KnowSys

Designing Twitter's Home Timeline

Ananya posts from a late train in Pune, and a few seconds later her friend Kabir opens the app in Bengaluru and sees it at the top of his feed, while a pop star with 100 million followers posts at the same moment. We'll design the system that gets both posts to the right screens, from one slow database query through fan-out on write, Redis timelines, the celebrity problem, Snowflake IDs, Manhattan, Earlybird search and the ranking pipeline, and look at why Twitter changed its mind more than once.

⏱ 55 min read◆ IntermediateAssumes: chapter 22 (Redis) helps, chapter 29 (partitioning) helps, chapter 38 (system design method) helps
Start reading

Ananya is on the evening train from Mumbai to Pune, and it has stopped in a field for the third time. She opens Twitter, types "40 minutes late again, and the chai has run out", and taps Post. Kabir, a friend from college, is on a sofa in Bengaluru. A few seconds later he opens the app, pulls down to refresh, and her post is the first thing on his screen.

A page of lined notepaper with a hand-drawn web page: a 'status' box with the word 'reading', a Set button, a list of statuses, and a 'watch?' link
Jack Dorsey's 2006 sketch of the idea that became Twitter. You set a short status, and other people can 'watch' you. Everything in this chapter is about the second half of the sketch: getting your status in front of everyone who watches you.Photo: Jack Dorsey, CC BY 2.0, via Wikimedia Commons

At the same moment, a pop star we'll call Nova posts a photo announcing her new album. Nova has 100 million followers, and Kabir is one of them. So the same refresh that shows Kabir his friend's complaint about the train has to show him Nova's post as well, in the right order, and so does the refresh of every one of Nova's other followers, most of whom have never heard of Ananya.

That sounds like a single lookup, "show Kabir the latest posts from the people he follows", and Twitter's first version did it that way. It's the version that gave Twitter its most famous error page. Why didn't it last? Look at the numbers. Ananya's post has to reach a few hundred people. Nova's has to reach a hundred million, and they all expect to see it within seconds. Meanwhile hundreds of thousands of people a second are opening the app, each with their own list of accounts they follow.

The question this chapter answers is: when someone posts, how does their post get onto the timelines of everyone who follows them within a few seconds, whether that's 300 people or 100 million? We'll start with the obvious design, find exactly where it breaks, and fix it step by step, the way Twitter did between 2008 and 2026: precomputed timelines in Redis, a special path for celebrities, IDs that sort by time, a database for the posts themselves, a search engine that turned out to matter for timelines too, and the ranking system that decides the order today.

01What we're building, and how big

1.1What it has to do

We'll call a short post a tweet, the name Twitter used until 2023 (X now calls them posts). A user's home timeline is the list of recent tweets from the accounts that user follows, newest first. It's what Kabir sees when he opens the app. Twitter also had a second kind of list, the user timeline: every tweet one account has sent, as shown on Ananya's profile page. Twitter's own engineers described the home timeline as a merge of the user timelines of everyone you follow, and that definition is a good one to hold on to, because every design in this chapter is a different way of doing that merge.

What we're designing comes down to a short list:

  1. Post a tweet, and store it so it's never lost.
  2. Deliver it to the home timeline of every follower, within a few seconds.
  3. Serve the home timeline fast whenever someone opens the app or scrolls.
  4. Search every tweet by its words, seconds after it's posted.
  5. Rank the timeline, so the posts a user most likely wants come first (Twitter added this in 2016; before then the timeline was strictly newest first).

And the qualities it needs:

  • Fast reads: opening the app should show a timeline in a fraction of a second.
  • Fresh: a tweet should reach followers' timelines within about five seconds, which was Twitter's own target in 2012.
  • Durable where it matters: a tweet, once posted, must never disappear. A timeline is a different matter, as we'll see: it can be rebuilt.
  • Available under spikes: the moments when everyone posts at once, a goal in a football final or a famous scene in a film, are exactly when people most want to read.

1.2How big is it?

The most detailed public numbers come from a talk called Timelines at Scale, given at QCon San Francisco in November 2012 by Raffi Krikorian, who then led Twitter's platform services group. At the time Twitter had about 150 million active users, who sent about 400 million tweets a day. That averages about 5,000 tweets a second, with a daily peak around 7,000 and more than 12,000 during big events.

What shapes everything else is on the other side. Krikorian put timeline reads at about 300,000 queries a second, against about 6,000 requests a second of writes. People read Twitter far more than they write to it. He described Twitter as "primarily a consumption mechanism, not a production mechanism".

In August 2013 Twitter's engineering blog gave newer figures: more than 500 million tweets on a typical day, about 5,700 a second on average, and a record one-second peak of 143,199 tweets a second, set on 3 August 2013 when viewers in Japan watching a TV broadcast of the film Castle in the Sky all posted at the same moment. That spike was about 25 times the normal rate. In 2023, the last time Twitter published a comparable figure, it was still about 500 million tweets a day.

Your turn: design it before reading on

In 2012 Twitter took in about 400 million tweets a day and made about 30 billion timeline deliveries a day (a delivery is one tweet landing on one follower's timeline). On average, how many timelines does one tweet land on? Now think about Nova. What does the average hide?

02Version 1: ask the database at read time

2.1The obvious design

Start with two tables in a relational database. One holds tweets: an ID, the author, the text and the time it was posted. A second holds follows: one row for each pair "Kabir follows Ananya". When Kabir opens the app, the server runs one query:

SQL
SELECT * FROM tweets
WHERE author_id IN (SELECT followee_id FROM follows WHERE follower_id = :kabir)
ORDER BY created_at DESC
LIMIT 20;

Read in plain words, it says: find everyone Kabir follows, find all their tweets, sort them by time and keep the newest twenty. This is the merge from section 1.1, done fresh every time someone looks. Doing the work when someone reads is called fan-out on read: the reader's request fans out to every account they follow. Twitter in 2006 and 2007 was a Ruby on Rails web application in front of a MySQL database, and it built timelines roughly this way.

Version 1: every timeline is a query
INSERTSELECTAnanya's phonepostsKabir's phonerefreshesRails appone request per process×manyMySQLtweets, follows
Step 1. Ananya posts. The app server inserts one row into the tweets table. Writing is cheap: one row, however many followers she has.
1 / 3

Writing is as cheap as it can be here: Ananya's tweet is one row, whether she has three followers or three million. All the cost lands on reads. If Kabir follows 300 accounts, his query has to look up the recent tweets of 300 authors and sort them together, and the database has to do it again for his next refresh, and for the next person's, 300,000 times a second. Krikorian's talk summed up this approach in one line: the naive version is "a massive select statement over all of Twitter", and "it was tried and died".

2.2The Fail Whale years

Slow reads weren't the only problem, and Twitter's 2013 account of these years lists the others. Tweets lived in a single MySQL database with one primary server, and when it filled up, Twitter started a new database for the newest tweets. That's called temporal sharding: splitting data by time, so each database holds one period. It seems tidy, but think about where the traffic goes. Every new tweet is written to the newest database, and most reads ask for recent tweets, which also live in the newest database. So the newest database is always the hottest, and adding more databases for older periods doesn't help at all.

Even the web servers had a ceiling. Each Rails server ran as a set of separate processes, each handling one request at a time, and Twitter measured them at 200 to 300 requests a second per machine. At that rate, keeping up with growth meant buying machines faster than anyone could install them.

When it all fell over, users saw an illustration by the designer Yiying Lu: eight small orange birds lifting a whale out of the sea in a net, captioned "Too many tweets! Please wait a moment and try again." Users called it the Fail Whale, and it became famous enough that the New York Times Magazine wrote about it in February 2009. In November 2013 Twitter's head of engineering told Wired that the company had taken it out of production that summer. It came to a head at the 2010 football World Cup. Twitter's engineering blog described how every shot on goal, penalty and card sent a flood of tweets that "repeatedly took its toll and made Twitter unavailable for short periods of time", while engineers worked through the nights for efficiency gains that the growth swallowed straight away.

Two children in orange Netherlands shirts watching the 2010 World Cup final on a television
Fans watching the 2010 World Cup final. Millions of people reacting to the same moment, at the same second, is the load pattern Twitter has to survive, and in 2010 it couldn't.Photo: Martin Thomas, CC BY 2.0, via Wikimedia Commons

2.3Off the monolith

After the World Cup, Twitter decided to rebuild instead of patching. Its 2013 write-up lists three goals: reduce the number of machines needed by ten times, isolate failures so one broken part couldn't take down the whole site, and let small teams ship features independently.

First came the servers themselves. Twitter already ran large services on the JVM, the Java virtual machine that runs Java and Scala programs: its search engine was written in Java, and its social graph store in Scala. A JVM server can handle many requests at once with threads, where each Rails process handled one. Twitter estimated a rewrite would give more than ten times the throughput on the same hardware, and by 2013 its JVM services were serving 10,000 to 20,000 requests a second per machine, against 200 to 300 for Rails.

Second, Twitter split the single Rails application, which Twitter called the monolith, into separate services, starting with what it called the "core nouns": a tweet service, a timeline service and a user service. Each service has one owner and one network interface, and teams agree on the interfaces, then build independently. For the calls between services, Twitter wrote Finagle, a library that gives every service the same connection pooling, retries, timeouts and load balancing, after early services each handled failures differently and set off storms of retries against each other.

From here on we follow the timeline service and the services around it. Its first problem was the one from 2.1: building every timeline from scratch on every read is far too slow. If reads outnumber writes fifty to one, we should do the work once, when the tweet is written, and make each read cheap.

03Fan-out on write: deliver every tweet ahead of time

3.1A mailbox per user

Here's the idea. Give every user a precomputed home timeline, like a mailbox, holding the IDs of the recent tweets they should see. When Ananya posts, look up her followers and drop her tweet's ID into each of their mailboxes. When Kabir opens the app, there's nothing to compute: read his mailbox, which already holds his timeline in order.

Doing the work when a tweet is written is called fan-out on write, because one write fans out into one insert per follower. It's the opposite of section 2's fan-out on read, and it turns the cost around: a read becomes one lookup, and a write becomes as many inserts as the author has followers. Since Twitter's reads outnumbered writes fifty to one, that trade was clearly worth making. By 2012 this was how Twitter built every home timeline, with the mailboxes kept in Redis, an in-memory key-value store whose values can be lists (chapter 22).

Ananya's tweet, fanned out on write (2012)
followers?insert IDread listAnanya's phoneWrite APIstores the tweetQueueFan-out serviceFlockwho follows whomTimeline cacheRedis, one list per user×thousandsTimeline serviceKabir's phone
Step 1. Ananya taps Post. The write API stores the tweet and puts a small job on a queue, then tells her phone it's done. She doesn't wait for delivery.
1 / 5

Notice the queue between the write API and the fan-out. Ananya's tap is answered as soon as the tweet is stored; delivery happens afterwards, in the background. In 2012 that was a necessity as much as a choice: the tweet-writing path was still Ruby, with only about 45 to 48 processes per machine, so each process had to hand the work off and hang up as fast as it could. It's also why Twitter's goal was phrased as delivery "in under 5 seconds" and not instantly. Fan-out takes time, and the design accepts a few seconds' delay in exchange for never making a reader wait.

3.2What's in the mailbox, and how much memory it takes

Each entry in a home timeline is tiny. Krikorian's talk listed what it held: the tweet's ID, the ID of the user who wrote it, and 4 bytes of flags marking whether it's a retweet, a reply or something else. An ID is a 64-bit number, so an entry is 8 + 8 + 4 = 20 bytes. Ananya's words aren't there at all.

0122.5tweet ID8bauthor ID8bflags4b64-bit Snowflake ID (section 6)retweet, reply, …who wrote it
One entry in a 2012 home timeline: 20 bytes. Ananya's tweet appears in Kabir's list as these three fields; the words she typed live elsewhere.

Each list was capped at 800 entries. Scroll back far enough and you hit the end. That cap is a design choice, and it's what makes the whole cache affordable:

Your turn: design it before reading on

Twitter kept a home timeline in memory for every active user: about 150 million of them in 2012, each with up to 800 entries of 20 bytes. How much memory is that? Twitter kept each timeline on three machines. What's the total?

Only active users got a timeline in memory. Twitter defined active as having logged in within the last 30 days, and the talk noted that the window could change depending on how much cache was available. If you hadn't logged in for a month, fan-out skipped you, and when you came back, Twitter rebuilt your timeline with a process it called reconstruction: ask the social graph who you follow, fetch each of their recent tweets from disk, merge them and load the result into Redis. That's version 1's fan-out on read, kept as the slow path for the rare user who needs it.

3.3Why a timeline can live in memory

Keeping something important only in memory sounds reckless, so it's worth being clear about why it's fine here. A home timeline doesn't record anything that isn't stored somewhere else: tweets are kept durably in their own store (section 7), and so is the follow graph. A home timeline is derived data: given who Kabir follows and what they've posted, it can always be computed again, as reconstruction does. Losing it costs time, not data.

Twitter still protected it. Each timeline was stored on three different machines in each data centre, each found through its own hash ring (chapter 29 explains how a hash ring maps keys to machines). A read asks whichever of the three answers fastest. If one machine dies, the other two still answer, so Twitter doesn't have to reconstruct every timeline that lived on it. Krikorian reported the timeline service's response time as about 5 milliseconds at the median and 100 at the 99th percentile (the time that 99 requests in 100 beat); only at the 99.9th percentile, when a timeline had to be rebuilt from disk, did it take a couple of hundred milliseconds.

04Inside one timeline: the data structure

4.1The operations a timeline needs

zoomTwitterTimeline cacheOne user's listHybrid list

Look at what fan-out and reads do to Kabir's list. Fan-out pushes a new ID onto the front. Reads take the first 20 or so from the front, then the next 20 as Kabir scrolls. And when the list grows past 800, the oldest entries fall off the back. In Redis terms that's LPUSH to add at the head, LTRIM to cut it back to 800, and LRANGE to read a slice. Chapter 22 covers Redis itself; here the question is what structure sits behind those three commands, because at 20 bytes an entry, the overhead of the structure can be bigger than the data.

Your first instinct is probably a linked list: each entry is a separate small block of memory holding the value plus a pointer to the next block, and usually one to the previous block as well. Pushing at the head and dropping from the tail are both cheap, because only a couple of pointers change.

Three boxes holding the numbers 12, 99 and 37, each with arrows pointing to the box before and the box after
A doubly linked list. Each value sits in its own node with two pointers, one to each neighbour. With 8-byte pointers, a 20-byte timeline entry carries 16 bytes of pointers, plus the memory allocator's own bookkeeping for every node.Image: Lasindi, public domain, via Wikimedia Commons

Its weakness is memory. Each node needs its two pointers, 16 bytes on a 64-bit machine, plus whatever the memory allocator adds to track each separate block. For values this small, the bookkeeping is about as large as the data, which would roughly double the terabytes from 3.2. Yao Yue, who worked on Twitter's cache team, made exactly this point in her 2014 talk on scaling Redis at Twitter: pointer overhead is high when each item is a small ID.

4.2One block, then many small blocks

At the other extreme, a structure can pack all the entries into one contiguous block of memory, one after another, with no pointers at all. Redis had such a structure, called a ziplist, and Twitter's original timeline design used ziplists exclusively. Memory-wise it's ideal: 800 entries cost about 800 × 20 bytes and almost nothing else.

Predict before you read on

Kabir's timeline is a single 16 KB ziplist holding 800 entries, newest first. Ananya's tweet arrives and must go at the front. What does Redis have to do?

So one structure wastes memory and the other wastes time. Twitter's fix was to combine them. A hybrid list is a linked list whose nodes are small ziplists. Each node holds a run of entries packed together, up to a fixed size in bytes, and the nodes are linked to each other with pointers. An insert at the head touches only the first small block, and if that block is full, a new block is linked in front of it. Trimming the oldest entries drops whole blocks from the tail. Now pointer overhead is paid once per block instead of once per entry.

Ananya's tweet arrives at Kabir's hybrid list (toy sizes: 3 entries per block, capped at 8)
Incomingfrom fan-outBlock 1head, newestBlock 2Block 3tail, oldestTrimmedpast the capDev95 s agoMeera5 min agoAnanya30 min agoDev1 h agoMeera2 h agoDev5 h agoMeera9 h agoAnanya2 days agoAnanyajust now
Step 1. Kabir's list holds 8 IDs in three blocks. Each block is a tiny ziplist: its entries sit side by side with no pointers between them. Only the blocks are linked.
1 / 5

Yue's 2014 talk gave the scale of this structure in one data centre: the timeline service's hybrid lists took about 40 TB of allocated memory, served about 30 million queries a second, and ran on more than 6,000 Redis instances. By 2017 Twitter's infrastructure blog named the cluster Haplo, "the primary cache for Tweet timelines", backed by a customised Redis implementing the hybrid list, written by the fan-out and timeline services and read by the timeline service. It handled 40 to 100 million Redis commands a second across the cluster, from about 800,000 requests a second to the service.

Decision

What should hold one user's timeline in memory?

Linked list
One node per entry, each with pointers to its neighbours.
  • Cheap insert at head and trim at tail
  • Pointers as big as the data
  • Many tiny allocations to manage
One packed block (ziplist)
All entries contiguous, no pointers.
  • Almost no overhead
  • Fast to scan
  • Head inserts copy the whole block
  • List length tied to a block-size limit
chosen
Hybrid list
A linked list of small packed blocks.
  • Head inserts touch one small block
  • Overhead paid per block
  • Trim drops whole blocks
  • A custom Redis to build and run

Twitter needed both properties at once: memory per entry close to the raw 20 bytes, because the cache was terabytes, and cheap inserts at the head, because every delivery is one. A hybrid list gets both by choosing the block size: big enough that pointers are a small fraction of each block, small enough that shifting one block is cheap. Its price was maintaining a private version of Redis.

05The celebrity problem

5.1When one tweet is 100 million inserts

Fan-out on write handles Ananya beautifully: 300 followers, 300 inserts, done in well under a second. Now Nova posts at the same moment.

Your turn: design it before reading on

Nova has 100 million followers. Suppose the fan-out service can deliver about 300,000 tweet IDs a second, roughly Twitter's whole delivery rate in 2012. How long until Nova's last follower has her tweet? What happens to Ananya's tweet while that's going on?

This wasn't hypothetical. In 2012 Lady Gaga had 31 million followers, and Katy Perry and Justin Bieber 28 million each. Krikorian's talk reported that fan-out to a million followers took about 3.5 seconds at the median, but the slowest 1% of large fan-outs took up to five minutes, and that it could take that long for a tweet from Lady Gaga to reach all her followers.

Delay also caused a stranger problem. A follower near the front of Lady Gaga's fan-out sees her tweet early and replies. The reply is an ordinary tweet with an ordinary fan-out, so it can reach some of her other followers before her original tweet does. Krikorian described exactly this: replies being seen before the tweets they reply to, a race condition that confused users. Sorting timelines by tweet ID couldn't fix it, because the original wasn't in the list yet.

5.2Measuring the trade

To see the trade in numbers, this program builds a small Twitter: 50,000 accounts, each following 50 others, chosen so that a few accounts are followed by almost everyone and most by almost nobody, the shape of real follower counts. Then it compares four policies. "Pure push" fans out every tweet. "Pure pull" fans out nothing and merges at read time, like version 1. In between, we fan out on write for everyone except accounts above a follower limit, and merge those accounts' tweets in at read time.

For each policy it reports the average number of timeline inserts per tweet, the worst single tweet (the most inserts any one tweet causes), the number of lists a read has to fetch, and the total work per tweet, counting each insert or list fetch as one unit and assuming 50 reads for every tweet, the 2012 ratio. random.choices with cum_weights picks followed accounts with probability proportional to 1/rank, so account 0 is the most popular. f & big is the set of big accounts one reader follows.

How much work fan-out costs under four policies
python
Python
import random
random.seed(65)
 
# A small Twitter: 50,000 accounts, each following 50 others.
# Popular accounts get picked far more often (account 0 most of all).
USERS, FOLLOWS = 50_000, 50
weights = [1 / (rank + 1) for rank in range(USERS)]
cum, total = [], 0.0
for w in weights:
    total += w
    cum.append(total)
 
followers = [0] * USERS
following = []
for user in range(USERS):
    picks = set(random.choices(range(USERS), cum_weights=cum, k=FOLLOWS))
    picks.discard(user)
    following.append(picks)
    for account in picks:
        followers[account] += 1
 
ordered = sorted(followers)
print(f"followers: median {ordered[USERS // 2]}, top account {ordered[-1]:,}")
 
# Every account tweets equally often; reads outnumber tweets 50 to 1
READS_PER_TWEET = 50
print(f"{'skip fan-out above':>18} | {'inserts/tweet':>13} | {'worst tweet':>11} | {'lists/read':>10} | {'work/tweet':>10}")
for limit in [None, 40_000, 5_000, 0]:
    if limit is None:
        big = set()
    else:
        big = {a for a in range(USERS) if followers[a] > limit}
    inserts = sum(followers[a] for a in range(USERS) if a not in big) / USERS
    lists = sum(1 + len(f & big) for f in following) / USERS
    worst = max((followers[a] for a in range(USERS) if a not in big), default=0)
    work = inserts + READS_PER_TWEET * lists
    label = "never (pure push)" if limit is None else f"{limit:,}"
    if limit == 0:
        label = "always (pure pull)"
    print(f"{label:>18} | {inserts:13.1f} | {worst:11,} | {lists:10.1f} | {work:10.0f}")
output
C++
followers: median 9, top account 49,492
skip fan-out above | inserts/tweet | worst tweet | lists/read | work/tweet
 never (pure push) |          42.4 |      49,492 |        1.0 |         92
            40,000 |          40.5 |      38,626 |        2.9 |        185
             5,000 |          30.9 |       4,896 |       12.5 |        655
always (pure pull) |           0.0 |           0 |       43.4 |       2169

Start with the first line: the median account has 9 followers and the top one has 49,492, nearly everyone. Now read the table top to bottom. Pure push does the least total work by far, 92 units per tweet against 2,169 for pure pull, because reads are fifty times more common and pure push makes every read a single list. That's the 2012 argument for fan-out on write, reproduced in a few lines. Pure pull makes each read fetch about 43 lists (the 50 follows, minus a few duplicates), and with fifty reads per tweet that dominates everything.

But look at the "worst tweet" column. Under pure push, one tweet from the top account costs 49,492 inserts, a thousand times the average of 42. That's Nova's tweet, and it's the column that hurts. Skipping only the two accounts above 40,000 followers costs a couple of extra list fetches per read and doubles the total work, but nobody's tweet is ever bigger than 38,626 inserts. Push the limit down to 5,000 and the worst tweet is under 5,000 inserts, at seven times pure push's total work. So the hybrid doesn't save work on average. What it buys is a bound on the worst single tweet. That bound decides how late Nova's post arrives and how long Ananya's waits behind it.

5.3The hybrid: merge celebrities in at read time

Twitter's answer, described in the 2012 talk as work in progress, was exactly this hybrid. For "high value users" with huge followings, stop fanning out. Keep their tweets only in their own user timelines, and when a follower reads, merge those few lists into the follower's precomputed timeline. Krikorian gave Taylor Swift as the example of an account whose tweets would be merged in at read time, and said that balancing the read and write paths this way saved "10s of percents of computational resources". Twitter never published the follower limit it used, or how many accounts were above it.

Kabir's refresh with the hybrid: one cached list plus a few merged ones
fan-outKabir's phoneTimeline servicemerges at read timeKabir's home listfanned out on writeSocial graphwhich big accounts?Nova's user timelinenot fanned outLeague's user timelinenot fanned outAnanya posts300 followers
Step 1. Ananya posts. She has 300 followers, well under the limit, so her tweet is fanned out on write as before and lands at the head of Kabir's home list.
1 / 5

That last step is a classic algorithm, the k-way merge. Each list is already sorted, newest first. Look at the head of each list, take whichever is newest, and advance that list. With a small structure called a heap, which always keeps the largest of its items on top, picking the newest of k heads takes about log₂ k steps, so merging three lists costs barely more than reading one.

Three rows of boxes, each row a short list of numbers in increasing order
Three sorted lists waiting to be merged. A k-way merge keeps one pointer per list, repeatedly takes the best head and advances that list. For Kabir, the rows are his home list, Nova's tweets and the league's tweets.Image: Amit6, public domain, via Wikimedia Commons

This program does Kabir's merge. It needs tweet IDs that sort by time, so it makes them the way Twitter does, with the Snowflake layout that section 6 explains: the time in milliseconds goes in the top bits, so a bigger ID means a newer tweet. Python's heapq.merge is a k-way merge over a heap, and reverse=True tells it the inputs are sorted largest first. islice takes the first five results.

Merge a cached timeline with two celebrities' tweets
python
Python
import heapq
from itertools import islice
 
EPOCH = 1288834974657                  # Snowflake's custom epoch, in ms
 
def snowflake(ms, machine, seq=0):     # 41 bits time | 10 bits machine | 12 bits sequence
    return ((ms - EPOCH) << 22) | (machine << 12) | seq
 
def seconds(tweet_id):                 # read the time back out of an ID
    return ((tweet_id >> 22) + EPOCH) / 1000
 
NOW = 1_791_658_800_000                # 10 Oct 2026, 19:00:00 UTC, in ms
author = {}
 
def tweet(who, seconds_ago, machine):
    tid = snowflake(NOW - seconds_ago * 1000, machine)
    author[tid] = who
    return tid
 
# Kabir's cached home timeline: IDs fanned out at write time, newest first
home = [tweet("Ananya", 2, 17), tweet("Dev", 95, 4), tweet("Meera", 300, 9),
        tweet("Ananya", 1800, 17)]
# Nova has 100M followers, so her tweets were not fanned out
nova = [tweet("Nova", 1, 31), tweet("Nova", 600, 22)]
# A second large account Kabir follows, also merged at read time
league = [tweet("CricketLeague", 45, 8)]
 
# Each list is already sorted newest first, so a k-way merge yields a sorted timeline
merged = heapq.merge(home, nova, league, reverse=True)
for tid in islice(merged, 5):
    print(f"{tid}  {author[tid]:<14} {NOW / 1000 - seconds(tid):6.0f} s ago")
output
C++
2108995977737269248  Nova                1 s ago
2108995973542907904  Ananya              2 s ago
2108995793187799040  CricketLeague      45 s ago
2108995583472582656  Dev                95 s ago
2108994723640283136  Meera             300 s ago

Kabir's merged timeline interleaves the three lists correctly: Nova's post from one second ago, then Ananya's from two seconds ago, then the league's, then the older entries from Kabir's cached list. Nothing in the program compares timestamps. It compares IDs, and the IDs are in time order because of how they're built. You can see it in the left column: each ID is smaller than the one above it. Notice also that heapq.merge stopped after five results without reading the rest of Kabir's list. A merge only reads as far as the page it's building, so merging a few lists at read time is cheap.

Decision

When should a tweet be delivered to its followers' timelines?

Fan-out on read
Build each timeline when it's read, from the authors' own lists.
  • Writes cost one insert
  • No per-user cache to keep
  • Every read merges hundreds of lists
  • Reads outnumber writes 50 to 1
Fan-out on write
Insert each tweet into every follower's precomputed list.
  • A read is one list lookup
  • Cheapest total work when reads dominate
  • One celebrity tweet is millions of inserts
  • Minutes of delay; replies before originals
chosen
Hybrid
Fan out on write, except for accounts above a follower limit, which are merged in at read time.
  • Bounds the worst single tweet
  • Reads still mostly one list, plus a few
  • Two delivery paths to build and test
  • The limit is a tuning knob

Twitter moved from pure fan-out on write to the hybrid around 2012 and 2013, because celebrities' fan-outs were delaying delivery for everyone and causing replies to show up before originals. It keeps fan-out on write's cheap reads for the great majority of tweets and caps the size of any single fan-out. As section 10 shows, this wasn't the last time Twitter changed its mind.

06IDs that sort by time: Snowflake

6.1Why the IDs matter

The merge in 5.3 only worked because tweet IDs sort by time. Fan-out sorts by them too, and so does every client that shows a timeline. In version 1 that came for free: MySQL handed out IDs from a counter that went up by one with each new row, so a newer tweet always had a bigger ID.

That stopped working once tweets were spread over many databases. Twitter's 2013 account explains the chain. To get away from the hot newest database of 2.2, Twitter used Gizzard, its own framework for sharding data across many MySQL servers, and built a tweet store called T-Bird on it: each new tweet's ID was hashed to pick a database, so writes and reads spread evenly over all of them. But a counter inside one database can't hand out IDs for the others, and two databases counting independently would hand out the same numbers. Twitter needed a new source of IDs, and in June 2010 it announced one called Snowflake.

6.2What Snowflake had to do

Twitter's 2010 announcement set out the requirements. Snowflake had to generate tens of thousands of IDs a second, with high availability, which ruled out a single central counter that every tweet had to wait for. The IDs had to be "roughly sortable": tweets posted around the same time should have IDs close together, because "this is how we and most Twitter clients sort tweets". Twitter aimed to keep IDs ordered to within a second. And the IDs had to fit in 64 bits, because changing the size of tweet IDs had already been painful once, with over 100,000 different programs built on Twitter's API.

Twitter rejected three alternatives. Ticket servers, a MySQL table that does nothing but hand out numbers, as Flickr used, didn't give the ordering without extra machinery. UUIDs, random identifiers that need no coordination, take 128 bits. ZooKeeper's sequential nodes, a counter kept by a coordination service, were too slow, and making every tweet wait on a coordinated service would lower availability.

Twitter's answer was to build each ID out of three parts: a timestamp, a worker number and a sequence number. Chapter 43 walks through the layout and what happens when a machine's clock jumps backwards (section 6, on deep dives), so here's the short version. The top 41 bits hold milliseconds since a custom starting moment (an epoch), 4 November 2010; the next 10 bits identify the machine (5 bits for the data centre and 5 for the worker, in the original code); and the last 12 bits count IDs within one millisecond on that machine.

A bar divided into three coloured sections labelled timestamp 41 bits, instance 10 bits and sequence 12 bits
A Snowflake ID: time in the top 41 bits, the machine in the next 10, a per-millisecond counter in the last 12. Because time is in the most significant bits, sorting IDs as plain numbers sorts them by time.Image: Sasmito Adibowo, CC BY-SA 3.0, via Wikimedia Commons

Because the time sits in the most significant bits, comparing two IDs as numbers compares their times first. The merge in 5.3 relied on exactly this. Two machines can't produce the same ID because their machine fields differ, and one machine can't repeat one because the time or the counter has moved on. Nobody has to talk to anybody to get an ID.

?Why only "roughly" sorted?

Because the machines' clocks don't agree exactly, and because the time is taken when the ID is created, not when the tweet becomes visible. A tweet created on a machine whose clock runs a few milliseconds fast can get a bigger ID than one created a moment later elsewhere. Twitter's announcement put it formally: the IDs are k-sorted, meaning every ID is at most a bounded distance from its sorted position, and Twitter aimed to keep that bound under a second. For a timeline, a few milliseconds of disorder between unrelated tweets is invisible.

6.3Reading a real ID

Since the time is inside the ID, you can read it back out of any tweet's ID with a shift and an add. This program decodes the ID of the selfie Ellen DeGeneres posted during the 2014 Academy Awards: shift right by 22 bits to drop the machine and sequence fields, add the epoch, and convert milliseconds to a date. & 0x1F keeps the low 5 bits after a shift, and & 0xFFF the low 12.

Decode the time and machine from a real tweet ID
python
Python
from datetime import datetime, timezone
 
EPOCH = 1288834974657   # Snowflake's epoch: 4 Nov 2010, 01:42:54.657 UTC
 
def decode(tweet_id):
    ms = (tweet_id >> 22) + EPOCH          # top 41 bits: milliseconds since the epoch
    datacenter = (tweet_id >> 17) & 0x1F   # 5 bits
    worker = (tweet_id >> 12) & 0x1F       # 5 bits
    sequence = tweet_id & 0xFFF            # low 12 bits
    when = datetime.fromtimestamp(ms / 1000, timezone.utc)
    return when, datacenter, worker, sequence
 
# The Oscars selfie, posted during the 2014 Academy Awards
when, dc, worker, seq = decode(440322224407314432)
print(f"created   {when:%Y-%m-%d %H:%M:%S.%f}"[:-3] + " UTC")
print(f"machine   datacenter {dc}, worker {worker}")
print(f"sequence  {seq}")
print(f"bits used {(440322224407314432).bit_length()} of 63")
output
C++
created   2014-03-03 03:06:13.747 UTC
machine   datacenter 1, worker 0
sequence  0
bits used 59 of 63

Line one gives 03:06 UTC on 3 March 2014, which is just after 7 pm on 2 March in Los Angeles, during the ceremony. Line two says which generator made the ID, and the sequence of 0 means it was the first ID that worker issued in that millisecond. And the last line shows there's room left: 41 bits of milliseconds last about 69 years from the 2010 epoch, so the top bit stays clear (signed 64-bit integers stay positive) until around 2080.

Those 64 bits caused one well-known problem. JavaScript stores numbers as 64-bit floating point, which holds integers exactly only up to 2⁵³, and Snowflake IDs are bigger. A web page that parsed a tweet ID as a number would silently round it to a different tweet. Twitter's API answered by sending every ID twice, as a number in id and as a string in id_str, and telling developers to use the string.

07Storing the tweets themselves: Manhattan

7.1Hydration

A timeline holds IDs, so before Kabir sees anything, each ID has to be turned back into a tweet: the text, the author's name and picture, the counts of likes and retweets. That step is called hydration. The timeline service takes the 20 IDs it's about to return and asks the tweet service for all of them at once, in parallel, and the user service for the authors. In 2012 the tweet service, Tweetypie, kept about the last six weeks of tweets in a memcached cluster (a plain in-memory key-value cache), and the user service, Gizmoduck, kept every user in cache. Krikorian's talk pointed out that the text of a tweet is "almost irrelevant" to most of the infrastructure: fan-out, timelines and ranking all work on IDs, and the words are only fetched at the end.

Behind the caches sits the store that holds every tweet ever posted, and that has to be durable. In 2012 it was T-Bird, the Gizzard-sharded MySQL from 6.1. From 2014 it was Manhattan.

7.2Manhattan

Twitter announced Manhattan in April 2014 as "our real-time, multi-tenant distributed database". Its motivation was operational: Twitter had many separate storage clusters, each built and run for one feature, and engineers spent too much time firefighting them and waiting for new capacity. Manhattan was built as one storage service for many teams, where an engineer could ask for space and throughput and start writing in seconds.

Its design choices follow Twitter's needs. It's a key-value store: you write a value under a key and read it back. It's eventually consistent by default, meaning a write is accepted by some replicas and reaches the others shortly afterwards, so a read can briefly return an older value. The announcement says most of Twitter's use cases "strongly favor availability over consistency", and that many of them couldn't accept even a few seconds of unavailability while a new primary server was elected. To keep replicas converging, Manhattan runs three repair mechanisms: a background process that continuously compares and reconciles replicas, read repair (fixing a stale replica when a read notices the difference), and hinted handoff (holding writes for a replica that's down and delivering them when it returns). For the cases that need it, strong consistency is available as an opt-in service built on a consensus algorithm (chapter 27), with compare-and-set operations within one data centre or across several.

Underneath, Manhattan offered three storage engines for different access patterns: seadb, a read-only format for data computed in batch jobs on Hadoop; a B-tree engine for read-heavy data; and an engine based on a log-structured merge tree (LSM tree) for write-heavy data. An LSM tree never updates data in place. It collects writes in memory, writes them out as sorted files, and merges files in the background, which turns many small random writes into a few large sequential ones. Chapter 18 covers how it works and what it costs.

Rows of small sorted files at level 0 being merged into fewer, larger sorted files at levels 1 and 2
Compaction in an LSM tree. New data is written as small sorted files, and background merges combine them into fewer, larger sorted files at each level. Writes stay sequential and fast, which suits a store that takes in thousands of tweets a second.Image: Ben Stopford, CC BY-SA 4.0, via Wikimedia Commons

By January 2017 Twitter's infrastructure blog described Manhattan as "the backend for Tweets, Direct Messages, Twitter accounts, and more", with read-only clusters handling tens of millions of queries a second and read-write clusters millions. Cassandra, which Twitter had adopted around 2010, had been deprecated in Manhattan's favour.

7.3How one tweet is laid out

zoomTwitterTweet storeManhattanKey layout

Twitter's 2023 open-source release of Tweetypie includes the code that writes tweets to Manhattan, and its comments document the layout. Each tweet is stored as many small key-value pairs under one primary key, the tweet ID, with a second-level key naming each piece:

C++
/[TweetId]/fields/internal/1          core fields, packed together
/[TweetId]/fields/internal/9 … /99    other built-in fields, one key each
/[TweetId]/fields/external/100 …      fields added later
/[TweetId]/metadata/delete_state      hard delete
/[TweetId]/metadata/soft_delete_state
/[TweetId]/metadata/scrubbed_fields/[FieldId]

A few details show the reasoning. The core fields that almost every read needs (the code says field IDs 2 to 8 of the tweet structure) are packed into one value under fields/internal/1, so a typical read is one lookup. Rarer fields get their own keys, so adding a feature means adding a key, with no change to existing tweets. Deletion is recorded as a state under metadata, separate from the content, so a delete can be undone or audited. And the tweet ID itself isn't stored inside the value, because it's already the key.

One more line in that file connects back to section 6. Manhattan orders keys as strings, and the string "10" sorts before "9". So the code pads both the tweet ID and the field ID with leading zeros to a fixed width, which makes string order match numeric order, and therefore time order, since the IDs are Snowflakes.

Decision

How should tweets be split across storage machines?

By time (to 2010)
Fill one database with the newest tweets; when it's full, start the next.
  • Simple
  • Old databases are quiet and read-only
  • All writes and most reads hit the newest database
  • Adding machines doesn't relieve the hot one
chosen
By hashed ID (T-Bird, then Manhattan)
Hash each tweet's ID to choose its shard.
  • Writes and reads spread evenly
  • Throughput grows with the number of shards
  • Needs IDs that don't come from one database's counter
  • A range of recent tweets is spread over every shard

Splitting by time put Twitter's whole write load on one primary. Hashing spreads it, and it's affordable because almost nothing reads "all tweets from the last minute" directly from the store: timelines come from the cache, and search from the search index. The one thing hashing broke, unique ordered IDs, is what Snowflake was built to fix. Chapter 29 compares the partitioning schemes in general.

8.1The opposite shape

Kabir sees Nova's announcement and searches for the album's name. Search has the opposite shape to the home timeline, and Krikorian's talk made the contrast explicit. For the home timeline, a tweet is written to every follower's list, and a read touches one list. For search, a tweet is written to exactly one search machine (plus its replicas), and a read asks every machine. Krikorian called search and pull "inverses" of each other.

Search works through an inverted index: for every word, a list of the tweets containing it. Chapter 62 explains how inverted indexes are built and queried, with postings lists, compression and ranking. Twitter's version had an extra demand that web search engines don't: a tweet had to be searchable within seconds of being posted, while the index served queries at the same time.

Twitter's real-time search engine is called Earlybird, introduced in October 2010 and described in a paper at the ICDE conference in 2012. It's built on Lucene, the open-source Java search library. As of autumn 2011, the paper reported over two billion queries a day, an average query latency of 50 milliseconds, and tweets typically searchable within 10 seconds of creation.

Searching for Nova's album: scatter, gather, merge
IN-MEMORY INDEX, HASH-PARTITIONEDone partitionNova's tweetIngesterstokenise, annotateEarlybird Apartition 1Earlybird Bpartition 2Earlybird Cpartition 3Root / Blenderfans the query outKabir searches
Step 1. Nova's tweet passes through the ingestion pipeline, which splits it into words and adds metadata such as the language.
1 / 4

8.2An index that's written while it's read

zoomTwitterSearchEarlybird serverActive segment

Earlybird's paper is unusually concrete about its data structures, and three choices stand out.

First, the index is cut into segments of a fixed number of tweets: 2²³, about 8.4 million, with 12 segments per server in 2011. New tweets fill one segment at a time, so at any moment only one segment is being written; the others are read-only, and once a segment is full it's rewritten into a compact, compressed form in the background.

Second, each posting, one occurrence of a word in a tweet, is a single 32-bit integer: 24 bits for the tweet's number within its segment and 8 bits for the word's position in the tweet. Eight bits are enough for a position because a tweet was at most 140 characters. A postings list is then just an array of integers. Searches want newest first, and the active segment appends new postings at the end, so a search reads each array backwards.

02432document ID within segment24bposition8bup to 2²⁴ tweetsword offset in a 140-character tweet
One Earlybird posting in 2011: the word appears in tweet number d of this segment, at position p. 24 bits cover the 2²⁴ tweets a segment could theoretically hold.

Third, the active segment has one writer thread and many reader threads, and they don't lock each other. The writer adds all of a tweet's postings, then increments a counter, maxDoc, holding the highest fully indexed tweet. That counter is declared volatile in Java, which places a memory barrier at the increment: everything the writer did before it is guaranteed to be visible to any thread that reads the counter afterwards (chapter 03 covers memory barriers). A reader starts each query by reading maxDoc and ignores any posting with a higher number, so it never sees half of a tweet.

One writer indexes Nova's tweet while Kabir's query reads the same segment
Index writerPostings arraysmaxDoc (volatile)Query threadappend 'album'read maxDoc = 4,000,000scan 'album' backwardsmaxDoc = 4,000,001next query reads 4,000,001
Step 1. The writer adds a posting for each word of Nova's tweet, number 4,000,001 in this segment, to the end of that word's array.
1 / 5

What does this buy? The paper measured it on 2011 hardware, a server with two quad-core processors and 72 GB of memory: about 17,000 queries a second on one 16-million-tweet segment, and indexing at 7,000 tweets a second under heavy query load, with the 99th-percentile latency under 200 milliseconds. Twitter's search had come a long way: in 2008 it ran on MySQL, inherited from Summize, the search start-up Twitter bought, and handled about 20 tweets and 200 queries a second.

8.3Earlybird in 2023

When Twitter published its recommendation code in 2023, the Earlybird README described the system's later shape. The index is split into three clusters: a realtime cluster holding public tweets from about the last seven days, a protected cluster for protected accounts' tweets over the same window, and an archive cluster holding every tweet ever posted, up to about two days ago. Ingesters read tweets from Kafka, a distributed log (chapter 23), and a root service fans each query out to the partitions and merges the answers.

That README holds a surprise for this chapter. It describes Earlybird's main use case as two things: search, and "Timeline In-network Tweet retrieval": given the list of accounts a user follows, find their recent tweets. By 2023, the search index had become the main source of tweets from the people you follow, for both the Following and For You tabs. That's fan-out on read, done by a search engine. To see why Twitter went back to reading at read time, we need to look at what the timeline had become: a ranked feed.

09Ranking: from newest-first to For You

9.1The funnel

Kabir's timeline in 2012 was every tweet from the accounts he follows, newest first. From 2016 Twitter ranked it, and by 2023 the default tab, For You, mixed tweets from accounts Kabir follows with tweets from accounts he doesn't. On 31 March 2023 Twitter published the code behind For You on GitHub as twitter/the-algorithm, with a blog post explaining it, and the published numbers are a good picture of what ranking costs.

Ranking has to narrow about 500 million tweets a day down to the few dozen Kabir will look at. Chapter 58, on Netflix's recommendations, explains the general shape of the answer: a funnel, where cheap methods cut millions of items down to a few thousand, and an expensive model scores only those few thousand. Twitter's funnel had three stages: fetch candidates, rank them, then filter and mix.

Candidate sources. For each request, the pipeline tries to pull the best 1,500 tweets out of hundreds of millions. On average half come from accounts Kabir follows, called in-network, and half from accounts he doesn't, called out-of-network. In-network candidates come from the search index, Earlybird, which ranks the followed accounts' recent tweets with a light model, a logistic regression (a simple weighted sum of features turned into a probability); the most important signal is Real Graph, a model that predicts how likely Kabir is to interact with each account he follows. Out-of-network candidates come from two kinds of source. One walks the graph of who engaged with what ("what did the people Kabir follows recently like?") on an in-memory engine called GraphJet; in 2023 that supplied about 15% of home timeline tweets. The other uses embeddings, vectors that place users and tweets so that similar ones are close: SimClusters groups accounts into 145,000 communities, recomputed every three weeks, and describes each user and tweet by the communities they belong to.

Ranking. All 1,500 candidates are then scored by the same model, the heavy ranker, a neural network of about 48 million parameters. Its input is about 6,000 features describing the tweet, Kabir and the relationship between them. Its output is ten probabilities, one for each way Kabir might engage: like, retweet, reply, open the author's profile, watch half a video, reply and get a reply back from the author, and so on, including negative ones such as "show less often" and reporting the tweet.

Filters and mixing. Finally, rules adjust the ranked list: remove tweets from accounts Kabir blocks or mutes, avoid too many consecutive tweets from one author, keep the in-network and out-of-network balance, require an out-of-network tweet to have a connection to someone Kabir follows, and thread replies with their originals. Then ads and account suggestions are mixed in, and the timeline is sent.

9.2Turning ten probabilities into one score

The heavy ranker's README in the companion repository, twitter/the-algorithm-ml, gives the formula that turns ten probabilities into one score: a weighted sum, each probability times a weight from a configuration file. It also lists the weights as of 5 April 2023:

Engagement predictedWeight
Like0.5
Retweet1.0
Reply13.5
Open profile and like or reply12.0
Watch at least half of a video0.005
Reply that the author engages with75.0
Click into the conversation and like or reply11.0
Click into the conversation and stay 2 minutes10.0
Negative feedback (show less, block, mute)−74.0
Report−369.0

Why would a report count 738 times as much as a like? The weights look extreme until you remember what they multiply. A like is far more likely than a reply that the author answers, so its probability is much bigger, and the README says the weights were originally set so that each term contributes about equally on average, then adjusted over time to move platform metrics. Read that way, the table seems to say the 2023 ranker prized conversations, especially ones the author joins, and punished tweets a user was likely to report.

All of this is expensive. Twitter's 2023 post says it ran about 5 billion times a day, finished in under 1.5 seconds on average, and needed 220 seconds of CPU time per run, nearly 150 times the time a user waits. That's possible because the 1,500 candidates are scored in parallel across many machines.

Kabir opens For You (2023): candidates, ranking, mixing
Kabir's phoneHome Mixerbuilds the feedEarlybirdin-network, ~50%Out-of-networkGraphJet, SimClustersHeavy ranker~48M params, 10 outputsTweet storehydrationFilters + mixingads, safety, diversity
Step 1. Kabir opens the For You tab. Home Mixer, the service that builds his feed, starts the pipeline.
1 / 6

9.3Why ranking undoes fan-out

Here's the connection to the rest of the chapter. Fan-out on write precomputes one thing: the newest 800 tweets from the accounts you follow, in time order. A ranked feed doesn't want that list. It wants a pool of recent candidates from everyone you follow, which it then reorders by predicted engagement, plus candidates from accounts you don't follow, which fan-out could never provide. Once ranking runs on every request anyway, the precomputed list saves less, and a source that can answer "recent tweets from these 400 accounts" quickly does the job. A real-time search index can do exactly that.

Twitter's 2023 post says it in one sentence: "We recently stopped using Fanout Service, a 12-year old service that was previously used to provide In-Network Tweets from a cache of Tweets for each user." The code shows the rest. Home Mixer's README lists the Following tab, the strictly reverse-chronological one, as built from an Earlybird candidate pipeline too. Twitter didn't publish a cost comparison or a full explanation for the switch, so the reasoning above is probably the closest we can get: it's what the published architecture implies.

10What changed after 2022

10.1What's documented

Elon Musk bought Twitter in October 2022, and in 2023 the company renamed itself X. Engineering publication mostly stopped, so much of what happened inside since then is unpublished: traffic figures, hardware, how timelines are stored. Three things are documented in public code and posts.

One is the 2023 release covered in section 9, which documented the retirement of the Fanout Service and the move of in-network retrieval to the search index.

Another came in January 2026, when X published a new repository, xai-org/x-algorithm, under the Apache 2.0 licence, with the code that builds the For You feed. It replaces much of the 2023 system. In-network posts now come from a service called Thunder, which, in the repository's words, "holds recent posts in memory as they are published, and returns those from the accounts a viewer follows". Its code shows the data structure: a concurrent hash map from each author's ID to a double-ended queue of that author's recent posts (post ID and creation time), fed from Kafka, with a cap per author and old posts trimmed after a retention window. A request passes in the viewer's following list, up to 10,000 accounts, and Thunder returns up to 1,200 recent posts from them, which are then filtered to the last 48 hours.

Look at the shape of that. It's a user timeline per author, kept in memory, and a home timeline assembled at read time by reading the user timelines of everyone the viewer follows: fan-out on read, the version 1 design, now affordable because everything lives in memory and the result is going to be ranked anyway. The 2012 hybrid treated a few celebrities this way; Thunder treats every author this way.

In-network posts in 2026: Thunder reads each author's recent posts at request time
following listAnanya postsNova posts100M followersPost eventsKafkaThunderauthor → recent posts×nHome MixerRustPhoenixtransformer rankerKabir's phone
Step 1. Ananya and Nova post. Each post is one event on a Kafka stream, whatever the author's follower count.
1 / 4

A third is the ranking model. X's 2026 repository describes Phoenix, a transformer, the kind of neural network behind today's language models, written in JAX (a Python library for numerical computing) with a Rust serving layer, used both to retrieve out-of-network candidates and to rank all of them. (The January release shipped a sample transformer ported from xAI's open Grok-1 model; an August 2026 update replaced it with the production code.) Phoenix reads the viewer's recent engagement history and predicts a probability for each of about 25 actions, positive (like, reply, repost, dwell time) and negative (mute, block, report), and, as in 2023, a weighted sum turns them into one score. Two of its documented design decisions matter for a system designer. When the model scores a batch of posts, each post is compared with the viewer's history but never with the other posts in the batch, so a post's score doesn't depend on what else happens to be scored alongside it, which makes scores consistent and cacheable. And embeddings are looked up by hashing, with no fixed vocabulary, so a post is representable the moment it's created.

10.2Reading the pendulum

Put the three eras side by side and Twitter's home timeline has swung from one end of the fan-out spectrum to the other and back:

YearsHow Kabir's in-network tweets were gatheredWhat drove it
2006–2010Fan-out on read: a database query per refreshSimplicity; it fell over under load
2010–2022Fan-out on write into Redis lists, with celebrities merged at read timeReads outnumbered writes 50 to 1, and a newest-first timeline could be precomputed
2023Fan-out on read from the Earlybird search indexRanking needs a candidate pool, not a fixed list
2026Fan-out on read from Thunder's in-memory per-author queuesThe same, with a store built only for this job

Each of these choices was probably right for the product and hardware of its time. Fan-out on write wins when the answer to a read can be computed in advance and reads dominate. Once the answer depends on a model run at read time, there's little left to precompute, and the cost of a hundred-million-insert fan-out buys nothing.

11The whole system

11.1Every box, and why it's there

Here's the design as it stood at its most documented point, around 2012 to 2017, with the hybrid fan-out, followed through one refresh.

Twitter's home timeline, end to end (hybrid fan-out era)
REBUILDABLE, IN MEMORYnew IDhydrateAnanya's phoneKabir's phoneWrite APITweetypieSnowflaketime-ordered IDsManhattanevery tweet, durableFan-outskips big accountsSocial graphFlockTimeline cachehybrid lists, ×3Timeline serviceread + mergeEarlybirdsearch index
Step 1. Ananya posts. The write API gets a Snowflake ID for the tweet, a 64-bit number with the time in its top bits.
1 / 6
ComponentWhat it doesAdded because
Timeline cacheOne capped list of tweet IDs per active user, in memoryBuilding timelines with a query on every read died at 300,000 reads a second (§2, §3)
Hybrid listA linked list of small packed blocks per timelineLinked lists waste memory; one packed block makes head inserts copy everything (§4)
Fan-out serviceInserts each tweet ID into followers' listsReads outnumber writes 50 to 1, so work belongs on the write (§3)
Read-time mergeMerges a few big accounts' tweets into the cached listOne celebrity tweet was millions of inserts and minutes of delay (§5)
SnowflakeTime-ordered 64-bit IDs without a central counterHash-sharded storage lost the database counter, and merges need IDs that sort by time (§6)
ManhattanDurable key-value store for every tweetTime-sharded MySQL put all load on one primary (§2, §7)
EarlybirdReal-time inverted index, partitioned by tweetSearch needs every word of every tweet, within seconds (§8)
Ranking pipelineCandidates, a heavy ranker, filtersThe feed became ranked, and later replaced fan-out for in-network tweets (§9, §10)

11.2From top to bottom

LevelThe choiceData structure or algorithm
SystemDo work on the side that happens less oftenPrecomputed per-user lists (2010–2022); per-author lists read at request time (2023–)
TimelineCap it and keep it in memory800 entries × 20 bytes; three replicas on three hash rings
One listPack entries, link blocksHybrid list: a linked list of ziplists
CelebritiesDon't fan them outk-way merge of sorted lists with a heap at read time
IDsSortable without coordinationSnowflake: 41 bits of milliseconds, 10 of machine, 12 of sequence
TweetsHash-partitioned key-value recordsManhattan keys /[TweetId]/fields/…, zero-padded so string order is numeric order
SearchPartition by tweet, scatter every querySegments of 2²³ tweets; 32-bit postings read backwards; one writer, volatile maxDoc
RankingFunnel from millions to 1,500 to a pageCandidate sources, a 48M-parameter network with 10 outputs, weighted sum

12What goes wrong, and what it cost

12.1Failures this design has to survive

What happensWhat the user seesWhat the design does
A timeline cache machine diesNothingReads go to one of the other two replicas; the lists are rebuilt later
A user returns after two monthsA slower first loadTheir timeline was evicted, so it's reconstructed from the social graph and the tweet store
A celebrity tweetsTheir tweet appears promptlyFan-out skips it, and each follower's read merges it in
A celebrity's tweet is fanned out anyway (pre-hybrid)Replies before the original; delays of minutesThe reason the hybrid exists
Everyone posts in the same second (143,199 a second in 2013)Nothing, if it worksWrites return after queueing; fan-out and indexing absorb the spike in the background
Two ID generators' clocks disagree slightlyNothing visibleIDs are only k-sorted; milliseconds of disorder don't matter in a timeline
A generator's clock jumps backwardsA short stall in new tweetsSnowflake refuses to issue IDs until its clock passes the last one used (chapter 43)
A query arrives mid-way through indexing a tweetNever half a tweetReaders ignore postings above maxDoc

12.2The tradeoffs, in one table

DecisionChosenGiven upWhy it was worth it
Where timeline work happens (2010–2022)On writeOne insert per follower per tweetReads outnumbered writes 50 to 1
CelebritiesMerge at read timeA slightly costlier readBounds the largest fan-out; no more minutes-long delays
Timeline storageMemory, capped at 800, three copiesDurability and history past 800It's a cache: the tweets and graph are stored elsewhere
List structureHybrid listA stock RedisNear-raw memory per entry with cheap head inserts
IDsSnowflake, k-sortedStrict ordering and 1-by-1 incrementsNo central counter; IDs sort by time
Tweet storageHashed, eventually consistent by defaultImmediate consistency everywhereAvailability, and no hot primary
Search writesOne partition per tweetEvery query asks every partitionCheap writes; queries are rarer than timeline reads
In-network retrieval (2023–)Fan-out on readPrecomputed listsA ranked feed wants candidates, not a fixed order

13Summary

  1. Twitter is read far more than it's written: in 2012, about 300,000 timeline reads a second against about 6,000 writes, and 400 million tweets a day.
  2. Building every timeline with a query on every read doesn't scale: each refresh merges hundreds of authors' tweets, and in Twitter's words, it "was tried and died".
  3. Fan-out on write moves the work to the rarer side: each tweet is inserted into every follower's precomputed list, so a read is one lookup.
  4. A home timeline is a cache, not a record: capped at 800 IDs of 20 bytes, kept in memory only for active users, replicated three times and rebuilt when lost.
  5. The hybrid list fits both needs of one timeline: packed blocks keep memory near the raw data, and linking small blocks keeps inserts at the head cheap.
  6. Celebrities break fan-out on write: 100 million followers means 100 million inserts and minutes of delay, so their tweets are merged in at read time with a k-way merge.
  7. Snowflake IDs sort by time without a central counter: milliseconds in the top 41 bits, machine and sequence below, so merging and sorting by ID sorts by time.
  8. Tweets themselves live in Manhattan, hash-partitioned and eventually consistent by default, laid out as many small keys under each tweet ID.
  9. Search is fan-out's inverse: one write per tweet into an in-memory Earlybird partition, and every query scattered to all partitions.
  10. Ranking turned the timeline into a funnel: about 1,500 candidates, a 48-million-parameter network predicting ten engagements, a weighted sum, then filters.
  11. The pendulum swung back: in 2023 Twitter retired its 12-year-old Fanout Service, and in 2026 in-network posts come from per-author lists read at request time.

14Build this

A timeline service with a celebrity limit.

  • Generate a follow graph like the one in 5.2, with 100,000 users and a heavy-tailed follower count. Write a Snowflake generator (41/10/12 bits) and use it for every tweet.
  • Run a local Redis. Implement fan-out on write: for each tweet, LPUSH its ID to every follower's home:{user} list and LTRIM it to 800, using a pipeline to batch the commands. Measure how long the top account's tweet takes to fan out.
  • Add a follower limit. Above it, only LPUSH to the author's user:{author} list. On read, merge home:{user} with the user lists of the big accounts they follow, using heapq.merge. Plot read latency and the largest fan-out time as you move the limit.
  • Make a reply race: have a follower near the start of a celebrity's fan-out reply immediately, and count how many followers see the reply before the original, with and without the limit.
  • Compare memory: store 800 IDs per user as a Redis list, then as a single packed string of 8-byte integers, and read the difference with MEMORY USAGE.

15Interview questions

beginnerWhat's the difference between fan-out on write and fan-out on read for a news feed?›

Fan-out on write delivers each new post to every follower's precomputed feed when it's written, so reading a feed is one lookup but writing costs one insert per follower. Fan-out on read stores each post once and builds the feed when someone opens it, by fetching and merging the recent posts of everyone they follow, so writes are cheap and reads are expensive. Which is better depends on the read:write ratio and the follower distribution. Twitter in 2012, with about 50 reads per write, chose fan-out on write for almost everyone.

beginnerWhy did Twitter store only tweet IDs in timelines, not the tweets themselves?›

Because a tweet appears in many timelines, on average about 75 in 2012 and up to tens of millions for celebrities. Storing the ID (with the author ID and a few flag bits, 20 bytes in all) keeps each timeline entry tiny, so a full 800-entry timeline is about 16 KB and every active user's timeline fits in a couple of terabytes of memory. The text is fetched once per read, in a parallel batch, from the tweet service and its cache; this is called hydration. It also means an edit or deletion only has to change one record.

intermediateA user with 100 million followers posts. Walk through what happens under pure fan-out on write, and how you'd fix it.›

Pure fan-out on write inserts the tweet's ID into 100 million lists. At a few hundred thousand inserts a second that takes minutes, so followers see the tweet at very different times, other users' tweets queue behind it, and replies from early recipients can reach later recipients before the original, which Twitter observed in 2012. The fix is a hybrid: don't fan out accounts above a follower limit; keep their tweets in their own user timeline, and at read time merge the few big accounts a user follows into that user's precomputed list. Because every list is sorted by time-ordered ID, a k-way merge with a heap does this in one pass.

intermediateWhy do Twitter's IDs need to be roughly sorted by time, and how does Snowflake do it without a central counter?›

Timelines, merges and clients sort tweets by ID, so a newer tweet needs a bigger ID. Once tweets were hash-sharded across many databases, no single auto-increment counter could issue IDs, and a central ID service would be a bottleneck and a single point of failure. Snowflake builds each 64-bit ID from milliseconds since a custom epoch in the top 41 bits, a machine ID in the next 10 and a per-millisecond sequence in the last 12. Machines never coordinate: their machine IDs differ, and each machine's time and sequence only move forward. Ordering is only k-sorted, within clock skew between machines, which is fine for a timeline.

deepWhy might a company that built fan-out on write later switch to fan-out on read?›

Fan-out on write precomputes a specific answer, a reverse-chronological list per user, and it pays off when that's what reads need and reads dominate. Twitter's 2023 recommendation post said it had stopped using its 12-year-old Fanout Service, and in-network tweets came from the Earlybird search index; by 2026 they came from Thunder, an in-memory store of each author's recent posts read at request time. Once every read runs a ranking model over a pool of candidates from all followed accounts, plus accounts the user doesn't follow, the precomputed order is thrown away, while the celebrity fan-out cost remains. Reading per-author lists at request time costs one append per post, regardless of followers, and gives the ranker exactly the candidates it wants.

deepHow does a real-time search index serve queries while it's being written to, without locks?›

Earlybird restricts writes to one thread per active segment and publishes progress through a single volatile counter, maxDoc. The writer appends all of a tweet's postings, then increments maxDoc; the volatile write is a memory barrier, so every earlier write is visible to any thread that reads the new value. A reader starts by reading maxDoc and ignores postings above it, so it may see some postings early but never acts on half a tweet. Postings are 32-bit integers in preallocated arrays, appended in time order and read backwards for newest-first results, and only one segment of about 8 million tweets is ever writable.

16Go deeper

check yourself
Twitter's 2012 timelines held 800 entries of 20 bytes for about 150 million active users. Roughly how much memory is one copy?›

800 × 20 bytes = 16 KB per user, times 150 million is about 2.4 TB. Twitter kept three copies, so about 7 TB before data-structure overhead, which is why the hybrid list's low overhead mattered.

Under the hybrid, a celebrity's tweet is posted. How many timeline inserts happen, and where does the work go?›

None into followers' timelines: the tweet goes only into the celebrity's own user timeline. The work moves to reads, where each follower's request merges the celebrity's recent tweets into their precomputed list.

Shift a Snowflake ID right by 22 bits and add 1288834974657. What do you get?›

The time the ID was created, in milliseconds since 1970 (Unix time). The 22 bits shifted away are the 10-bit machine field and the 12-bit sequence.

Why does Earlybird store postings in time order and read them backwards?›

New tweets get higher document numbers and are appended at the end of each postings array, which is cheap. Real-time search wants newest first, so readers iterate from the end, and every array position is a valid place to start even while the writer keeps appending.

Raffi Krikorian, 'Timelines at Scale' (QCon San Francisco, 2012)

The source of most 2012 numbers here: 300K reads a second, 800-entry Redis timelines, fan-out times and the celebrity problem. Video on InfoQ; a detailed written summary is on High Scalability (July 2013).

'New Tweets per second record, and how!' (Twitter Engineering, August 2013)

The 143,199-tweets-a-second record, the 2010 World Cup, and the move from the Rails monolith to JVM services, Finagle, T-Bird and Snowflake.

'Announcing Snowflake' (Twitter Engineering, June 2010)

The requirements, the rejected alternatives and the k-sorted guarantee. The 2010 code is in twitter-archive/snowflake on GitHub.

'Manhattan, our real-time, multi-tenant distributed database' (Twitter Engineering, April 2014)

Why Twitter built one storage service for many teams, its consistency model and its three storage engines.

Busch et al., 'Earlybird: Real-Time Search at Twitter' (ICDE 2012)

A six-page paper with the segment layout, 32-bit postings, slice pools and the single-writer memory-barrier design.

'Twitter's Recommendation Algorithm' (Twitter Engineering, March 2023) and twitter/the-algorithm

The For You pipeline, candidate sources, the heavy ranker and its weights, and the retirement of the Fanout Service.

xai-org/x-algorithm (X, 2026)

The current For You code under Apache 2.0: Thunder's in-memory per-author post store, the Phoenix transformer ranker and the filters.

Yao Yue, 'Scaling Redis at Twitter' (2014) and 'The Infrastructure Behind Twitter: Scale' (2017)

The hybrid list and the timeline cache's size, and the Haplo cluster in Twitter's own infrastructure overview.

System Design Method

Fan-out on write versus read as an access-pattern decision, using the same Twitter numbers. Chapter 38.

Redis

Listpacks, quicklists and what one key costs, the machinery under section 4. Chapter 22.

Interview Systems

Snowflake's bit layout and what happens when the clock goes backwards. Chapter 43.

Search Engines

Inverted indexes, postings compression, ranking and sharding, the general version of Earlybird. Chapter 62.

Designing Netflix Recommendations

The candidate-and-ranker funnel, embeddings and experiments behind a ranked feed. Chapter 58.

Partitioning

Hash rings, and why hashing beats time ranges for write-heavy data. Chapter 29.

Storage Engines

LSM trees and B-trees, the engines inside Manhattan. Chapter 18.