Problem & requirements
Users publish posts and see a feed of their friends' posts. Clarify whether it is mobile, web, or both; how the feed is ordered (assume reverse chronological for simplicity, though real systems rank); how many friends a user can have (assume up to 5,000); the daily active user count (assume 10 million); and whether posts contain images and videos (assume yes).
There are two flows to design. Feed publishing writes a post and propagates it to friends' feeds. News feed building assembles a user's feed when they open the app. Everything else is optimisation of these two paths.
Back-of-the-envelope estimation
Assume 10M DAU, each posting twice a day: 20M posts per day, or 20,000,000 / 86,400 ≈ 230 posts per second. If the average user has 300 friends, fan-out on write performs 20M × 300 = 6 billion feed insertions a day, about 70,000 per second. That multiplication is why fan-out strategy matters so much.
If each feed cache entry stores a pair of 8-byte IDs (post ID and author ID, 16 bytes) and we keep the latest 500 entries per user, that is 8 KB per user and 10M × 8 KB = 80 GB for all active users, which fits easily in a Redis cluster. Double-check: 500 × 16 B = 8,000 B; 8,000 × 10^7 = 8 × 10^10 B = 80 GB.
APIs and high-level design
Two HTTP APIs cover the core. POST /v1/me/feed publishes a post with content and an auth_token. GET /v1/me/feed retrieves the caller's feed. Both go through a load balancer to web servers that authenticate and rate limit.
On publish, web servers call the post service, which stores the post in the post database and cache; the fan-out service, which pushes the new post into friends' feeds; and the notification service, which alerts friends. On read, the news feed service fetches the list of post IDs from the news feed cache and turns it into fully populated posts.
Fan-out on write vs fan-out on read
Fan-out on write (push) precomputes feeds: when a post is published, its ID is inserted into every friend's feed cache immediately. Reads are then very fast because the feed is already assembled. The downsides are that a user with millions of followers triggers millions of writes, the hotkey or celebrity problem, and that work is wasted on inactive users who may never open the app.
Fan-out on read (pull) builds the feed only when a user asks for it, by fetching recent posts from everyone they follow and merging. Nothing is wasted on inactive users and celebrity posts cost nothing extra at write time. But reads become slow, since one request may touch hundreds of timelines.
The practical answer is a hybrid: push for ordinary users, so most feeds are precomputed, and pull for celebrities, whose posts are merged in at read time. Consistent hashing helps spread hot requests and data across cache nodes.
The fan-out service in detail
When a post arrives, the fan-out service first gets the author's friend IDs from a graph database, which is well suited to relationship queries and friend recommendations. It then fetches friends' info from the user cache and filters by settings: someone who muted the author, or a post shared only with selected friends, should not reach every follower.
The filtered list and post ID go onto a message queue, and fan-out workers consume them and write <post_id, user_id> entries into the news feed cache. Only IDs are stored, not full posts or user objects, which keeps memory low. Each feed is capped to a few hundred entries because almost nobody scrolls through thousands of old posts; requests for older items can fall back to the database.
post_id into each feed list. Celebrity authors skip this path and are merged in at read time.Feed retrieval and cache architecture
To build a feed, the client calls GET /v1/me/feed. The news feed service reads the list of post IDs from the feed cache, then hydrates them by fetching usernames, profile pictures, and post contents from the user and post caches, and returns JSON. Media such as images and videos is stored in a CDN so it loads quickly from a nearby edge.
Caching is so central that it is worth splitting into tiers:
- News feed: lists of post IDs per user.
- Content: post data, with very popular posts kept in a separate hot cache.
- Social graph: follower and following relationships.
- Action: whether a user liked, replied to, or shared a post.
- Counters: like, reply, follower, and following counts.
Trade-offs & alternatives
Every choice here trades write cost against read latency and freshness. Push gives fast reads but expensive writes and wasted effort; pull is the reverse. The hybrid adds complexity at read time because two sources must be merged and sorted. Storing IDs rather than objects saves memory but requires a hydration step. Reverse chronological ordering is simple and predictable, while ranked feeds need a scoring service and features from the action and counter caches.
For storage, the post data can live in a relational or NoSQL database with read replicas, while the social graph benefits from a dedicated graph store or a TAO-like association cache on top of a sharded database.
Scaling, failure handling & wrap-up
Scale the web tier horizontally and keep it stateless. Scale databases with vertical then horizontal scaling, primary-replica replication for reads, and sharding. Use the message queue to decouple publishing from fan-out so a burst of posts does not slow the publish API, and so failed fan-out jobs can be retried. If the feed cache loses data, feeds can be rebuilt by pulling from recent posts, at a temporary latency cost.
Monitor QPS during peak hours and feed refresh latency. A strong answer explains the fan-out trade-off clearly, chooses a hybrid with a defined celebrity threshold, and shows that the feed cache stores only IDs.
Key numbers
Key terms
- Fan-out on write
- Precomputing feeds by pushing each new post into every follower's feed at publish time.
- Fan-out on read
- Assembling a feed at request time by pulling recent posts from everyone the user follows.
- Hotkey / celebrity problem
- A single account with so many followers that pushing its posts overwhelms the system.
- Hydration
- Turning a list of IDs into full objects by looking up their data in caches or databases.
- Graph database
- A store optimised for traversing relationships such as friendships and follows.
- News feed cache
- A per-user list of recent post IDs that makes feed reads fast.
Common mistakes
- Choosing pure fan-out on write without addressing users with millions of followers.
- Storing full post objects in each user's feed cache and multiplying memory use.
- Ignoring privacy settings and mutes when fanning out posts.
- Fanning out synchronously in the publish request instead of via a queue.
- Serving images and videos from application servers rather than a CDN.
Further study
- TAO: Facebook's Distributed Data Store for the Social Graph (2013)
- Twitter timelines at scale (Raffi Krikorian, QCon 2013)
- Feeding Frenzy: Selectively Materializing Users' Event Feeds (SIGMOD 2010)
- Scaling Memcache at Facebook (NSDI 2013)