System Design Problem

Design a Timeline and Tweet Service (Twitter)

Commonly Asked By:TwitterMetaLinkedInUber

Interview Setup

Interview Prompt

Design Twitter's home timeline. Users follow other users and receive a personalized, ranked reverse-chronological feed of tweets. Support retweets, list-based feeds, and handle accounts with tens of millions of followers without latency degradation.

Clarifying Questions (ask before designing)

QuestionWhy it matters
Are we designing the home timeline feed only, or also the user profile timeline and search indexing?The home timeline requires fan-out aggregation across hundreds of followed accounts, whereas a user profile timeline is a direct query against a single author partition in Cassandra.
Is the feed strictly reverse-chronological or ranked algorithmically?Algorithmic scoring introduces a machine-learning ranking stage, transforming the read path from a simple sorted merge into a multi-stage scoring, deduplication, and filtering pipeline.
What follower threshold differentiates normal accounts from high-follower celebrity accounts?This boundary dictates whether an incoming tweet is pushed into follower timeline caches in Redis or pulled on demand at query time.
How do native retweets and quote tweets differ in storage, fan-out, and deduplication?Native retweets reference an existing tweet entity via retweet_of and require deduplication during timeline merge, whereas quote tweets create a new independent tweet with custom commentary referencing quoted_tweet_id.
How do custom list feeds differ from the main home timeline?List feeds require distinct timeline caches in Redis keyed by list ID, where fan-out targets subscribed list timelines based on a reverse creator-to-list membership index rather than user follower lists.

Scope

In scope

  • Home timeline generation via hybrid push and pull fan-out architecture with Redis serving cache
  • Cassandra CommitLog CDC event streaming bridge to Kafka to mitigate dual-write loss
  • Retweet, quote tweet, conversation thread table indexing, and list-based feed mechanics
  • 7-stage timeline read pipeline: Candidate Generation, Visibility Filtering, Celebrity Merge, Retweet Deduplication, Ranking, Hydration, and Response
  • Follow, refollow, and unfollow lifecycle handling

Out of scope (state explicitly)

  • Full-text search engine inverted index infrastructure (handled by dedicated search clusters)
  • Ad auction placement and monetization pipelines
  • Direct messaging infrastructure
  • Complex trending topic velocity calculation (delegated to stream-processing systems)

Functional Requirements

Architectural Context: Focus here on Twitter-specific constraints including 280-character tweets, hybrid push-pull fan-out thresholds, retweet and quote deduplication, list-based feeds, conversation threads, and Cassandra CommitLog CDC event streaming. Review the News Feed System for the foundational fan-out pattern.

  • User Timeline: View a user's own tweets and retweets in reverse-chronological order on their profile page.
  • Home Timeline (Feed): Aggregate, filter, rank, and render tweets from all followed accounts with sub-second retrieval.
  • Follow and Unfollow Graph: Follow or unfollow creators with bounded backfill on follow and lazy deletion on unfollow.
  • Interactions: Like, native retweet, reply, and quote-tweet with real-time engagement counts, deduplication, and cascade deletion tombstones.
  • Conversation Threads: Cluster hierarchical replies under a dedicated conversation query table for single-partition thread metadata retrieval followed by tweet hydration.
  • List Feeds: Create and view custom list feeds with dedicated timeline caching and creator-to-list membership routing.
  • Tweet Search: Ingest and index new tweets in real time for hashtag and keyword discovery.
  • Trending Topics: Publish tweet stream events to downstream streaming systems (detailed in Trending Topics).
  • Notifications: Asynchronously alert users for mentions, likes, retweets, and new followers without blocking the tweet creation flow.

Non-Functional Requirements

Twitter systems prioritize read latency, horizontal scalability, event publication durability, and continuous availability across global audience surges.

  • Low Read Latency: Serve home timeline requests in under 200 ms at p99.
  • High Availability: Maintain 99.99% availability with graceful fallback to reverse-chronological feeds during ranking service degradations.
  • Massive Scalability: Support 400 million daily active users and 500 million tweets per day.
  • Eventual Consistency: Accept brief propagation delays of up to 5 seconds for fan-out timeline updates.
  • Read-Heavy Architecture: Optimize for a 20:1 read-to-write traffic ratio (10 billion timeline views vs 500 million tweets published daily).
  • Dual-Write Mitigation: Remove application-level database-to-Kafka dual-write windows by bridging committed Cassandra mutations to Kafka via CommitLog CDC.

Capacity Estimations

Evaluating daily tweet write volume against follower fan-out amplification dictates the memory footprint of Redis timeline caches and database storage growth.

MetricCalculationValue
Daily active users (DAU)Given product scale assumption400M
Tweets created dailyEstimated ~1.25 tweets per user daily500M
Average tweet write throughput500M tweets ÷ 86,400 seconds~5,787 writes / sec (rounded to 6K for planning, peak ~30K)
Home timeline reads daily400M DAU x 25 views per user10B timeline reads / day
Average timeline read throughput10B reads ÷ 86,400 seconds~115,740 reads / sec (peak ~500K)
Read-to-write traffic ratio10B reads ÷ 500M writes20:1
Average tweet metadata payloadText, author metadata, and media references1 KB
Daily tweet text storage500M x 1 KB500 GB / day (182.5 TB / year)
Daily media storage volume50M media tweets (~10% of 500M daily tweets) x 2 MB average100 TB / day (Amazon S3)
Fan-out write volume (avg 200 followers)500M x 200 followers100B fan-out writes / day (~1.2M writes/sec)

Throughput and Capacity Analysis

  • Write Throughput: 500 million tweets daily generates an average of ~5,787 writes per second (often rounded to 6,000 writes/sec for capacity planning), with peak surges reaching 30,000 writes per second during major global events.
  • Read Throughput: 10 billion timeline views daily translates to approximately 115,740 reads per second on average and up to 500,000 reads per second at peak.
  • Fan-Out Write Amplification: Fanning out 500 million tweets across an average of 200 followers produces roughly 100 billion cache write operations daily (~1.2 million Redis writes/sec).
  • Storage Growth: 500 GB of text metadata daily totals 182.5 TB per year, while media uploads (~10% of tweets containing media averaging 2 MB) generate 100 TB daily managed in Amazon S3.

Architecture Diagram

Interview strategy: Establish the celebrity fan-out threshold before designing data pipelines. A single misclassified high-follower account can generate tens of millions of redundant writes per tweet. Emphasize Cassandra CommitLog Change Data Capture (CDC) to bridge durable database commits to Kafka, removing the application-level dual-write failure window.

The system implements a hybrid fan-out architecture: normal users push tweets into follower caches in Redis at write time, while high-follower celebrity accounts are pulled on demand at read time.

Loading...

Hybrid Fan-Out Strategy

The core fan-out trade-off matches principles from the News Feed System. Twitter introduces an adaptive celebrity threshold to prevent linear follower-proportional write amplification while maintaining instant timeline reads.

User TypeFollowersStrategyReason
Normal creators (99.9%)< 10KFan-out on Write (Push)Pre-computes home timeline at write time so reads are fast bounded cache lookups
Celebrity creators (0.1%)≥ 10KFan-out on Read (Pull)Pulls recent tweets at query time to avoid enqueuing millions of writes per post

Step-by-Step Delivery Lifecycle

  1. A creator publishes a tweet, and the Tweet Service validates the content, writes the tweet record to the authoritative Cassandra store, and returns HTTP 201 Created to the client.
  2. A Cassandra CDC connector tails committed mutation logs on disk and publishes a message to the tweet-events Kafka topic with connector-level at-least-once delivery guarantees.
  3. The Fan-Out Service consumes the event and queries the Social Graph Service for the creator's follower count.
  4. If the creator has fewer than 10,000 followers, the service retrieves the follower list and executes pipelined ZADD writes of the tweet_id into each follower's Redis timeline sorted set.
  5. If the creator has 10,000 or more followers, the service skips push fan-out and records the tweet in the author's recent tweet serving cache in Redis only.
  6. When a follower requests their home timeline, the Timeline Service executes the 7-stage read pipeline: fetches a bounded window of ~200 precomputed candidates from Redis, pulls recent tweets from followed celebrities, filters deleted or blocked content, deduplicates retweets, scores candidates via the ML Ranking Service, hydrates the top 20 items, and returns the response.

Component Deep Dives

Comprehensive breakdown of tweet ingestion, Cassandra CDC event streaming, 7-stage timeline assembly, interaction and counter pipelines, list feeds, social graph indexing, machine-learning ranking, and asynchronous stream processing.

Tweet Service & Cassandra CommitLog CDC Pattern

The Tweet Service handles tweet validation and durable persistence, ensuring write API responses complete in under 100 ms while eliminating application-level dual-write failure windows.

When a user posts a tweet, writing directly to both the database and Kafka creates a dual-write failure risk where the database write succeeds but Kafka publication fails, resulting in ghost tweets that never fan out. To maintain consistency without requiring cross-partition distributed transactions, the Tweet Service commits the tweet record directly to the authoritative Cassandra tweets table.

Cassandra writes mutations sequentially to its on-disk CommitLog before acknowledging the write. A dedicated Cassandra Change Data Capture (CDC) connector acts as a durable, replayable bridge tailing the commit log on disk and publishing events to the tweet-events Kafka topic with connector-level at-least-once delivery, subject to connector uptime, offset tracking, and commit-log retention. The Tweet API responds with HTTP 201 Created as soon as Cassandra commits, leaving fan-out, search indexing, and push notifications to asynchronous downstream workers.

Interaction Service & Real-Time Counter Pipeline

User interactions including likes and retweets require low-latency serving counters alongside durable event logging.

When a user interacts with a tweet (POST /api/v1/tweets/{tweet_id}/like), the Interaction Service records the state durably in relational storage, publishes a like-events message to Kafka, and atomically updates real-time counters in Redis (HINCRBY tweet:counters:{tweet_id} like_count 1). During Step 6 (Hydration) of the timeline read pipeline, these counters are merged into the tweet object payload, while the ML Ranking Service consumes real-time interaction velocities to score and surface trending candidate tweets.

Timeline Service & 7-Stage Read Pipeline

The Timeline Service orchestrates timeline assembly across precomputed Redis caches and live candidate streams.

While Redis retains up to 800 tweet IDs per active user (subject to a sliding 48-hour TTL) to support continuous scrolling, each home timeline request evaluates a bounded working set through a structured 7-stage pipeline:

  1. Candidate Generation: Fetches a bounded window of ~200 precomputed tweet IDs from the user's Redis timeline sorted set (timeline:home:{user_id}) for normal followees, and queries the latest ~100 recent tweet IDs from followed celebrity caches (timeline:user:{celeb_id}) in parallel, assembling roughly 300 candidate tweet IDs.
  2. Visibility Filtering: Checks cached follow graph sets, mute lists, and block lists in Redis to filter out content from users no longer accessible.
  3. Celebrity Merge: Merges the precomputed candidate list with the celebrity tweet streams using an in-memory multi-way merge ($O(C \log S)$ where $C$ is the total candidate items consumed, roughly 300 items, and $S$ is the number of sorted candidate streams).
  4. Retweet Deduplication: Identifies duplicate entries where a user received both the original tweet and a retweet of that same tweet, collapsing them into a single entry displaying a retweet attribution badge.
  5. Ranking: Passes the consolidated candidate pool through the ML Ranking Service to score tweets by engagement velocity and author affinity (or sorts strictly by timestamp if the user selected the chronological "Latest Tweets" view).
  6. Hydration: Extracts only the top 20 candidate tweet IDs for the requested page and hydrates full tweet metadata, media URLs, real-time like/retweet counts, and author profiles from the Redis Tweet Cache (falling back to authoritative Cassandra records on cache miss).
  7. Response Assembly: Attaches pagination cursors and returns the final payload to the client. Under normal healthy warm-cache conditions with local network hops, this in-process pipeline executes within an illustrative ~12 to 15 ms internal serving budget, comfortably meeting our formal production SLO of p99 < 200 ms.

Fan-Out Service

The Fan-Out Service operates as a scalable Kafka consumer group dedicated to timeline distribution.

For each incoming tweet, the worker checks the author's follower count against the 10,000 threshold. Below the threshold, it fetches the follower list from the Social Graph Service and performs pipelined ZADDoperations into each follower's Redis timeline sorted set. For accounts at or above the threshold, push fan-out is skipped, and the tweet ID is recorded in the author's recent tweet cache in Redis.

Because Kafka provides at-least-once delivery, fan-out workers must be idempotent. Redis ZADD is naturally idempotent when scoring by tweet ID timestamp, ensuring network retries do not produce duplicate entries. The same worker fleet handles deletions (executing ZREM across follower timelines) and retweet distributions.

List Feed Architecture

List feeds represent curated subsets of creators and follow a distinct fan-out model separate from follower feeds.

The architecture explicitly separates list members (creators in the list) from list subscribers (users following the list). The system maintains a single precomputed timeline sorted set in Redis keyed by list ID (timeline:list:{list_id}). When a creator tweets, the Fan-Out Service queries a reverse membership index (mapping creator ID to list IDs containing that creator) and pushes the tweet ID into those list timeline caches.

Because the list timeline is precomputed once per list, having millions of list subscribers generates zero subscriber-proportional write amplification on tweet creation. Fan-out occurs once per list containing the creator rather than once per subscriber, and subscribers simply read that single precomputed list cache. If an individual creator is an extreme outlier included in tens of thousands of lists, the system uses the creator-to-list reverse membership index to pull that creator's tweets on read into list views rather than pushing to every list cache.

Social Graph Service & Follow Lifecycle

The Social Graph Service manages relationships and handles follow, unfollow, and refollow transitions.

  • Follow Event: When User A follows User B, the service executes a bounded backfill, inserting the latest 20 to 50 tweets from User B into User A's home timeline cache in Redis.
  • Refollow Event: If User A unfollowed and later refollows User B, the same bounded backfill executes rather than resurrecting arbitrarily stale historical entries.
  • Unfollow Event: Handled via lazy deletion. The user's follow set is cached in Redis with a short TTL. When the timeline is assembled at read time, tweets from unfollowed accounts are filtered out during candidate evaluation, while a background compaction job cleans up aged entries.
  • Block Event: Triggers an immediate read-time filter and enqueues an asynchronous task to purge the blocked user's tweets from the cache.

Search Service and Elasticsearch

Near-real-time search indexing keeps recent tweets discoverable within seconds of publication.

The Search Indexer consumes tweet-events asynchronously and indexes documents into an Elasticsearch cluster. Each indexed document includes the tweet ID, author ID, content, hashtags, mentions, creation timestamp, and initial engagement signals. At query time, BM25 text relevance is combined with recency decay and engagement boosts to rank search results. Full inverted index architectures are explored in Search Engine.

Trending Topics

Velocity-based hashtag tracking is decoupled into a dedicated stream-processing pipeline.

The core timeline architecture publishes raw tweet events to Kafka, allowing a downstream streaming cluster (such as Apache Flink using sliding windows and Count-Min Sketch) to compute trending hashtags per geographical region. The complete streaming architecture is covered in Trending Topics.

Notification Service

Notifications are fully asynchronous and never sit on the critical path of tweet creation.

A dedicated consumer group listens to tweet-events, like-events, and follow-eventsto dispatch mobile push alerts via Apple APNs and Google FCM. Similar events are aggregated into digest notifications (such as "Alice and 15 others liked your tweet") to prevent notification storms during viral surges.

ML Ranking Service

The home timeline applies machine-learning models to rank candidate tweets for user engagement.

After candidate generation produces roughly 300 candidate tweet IDs, a lightweight gradient-boosted decision tree scores each tweet by recency decay, author affinity, engagement velocity (likes, retweets, replies), content media type, and verification status. The scoring stage executes in under 10 ms for a candidate batch. If the ranking service experiences degradation, the Timeline Service falls back to chronological sorting.

Analytics Pipeline

Analytics pipelines operate on independent SLAs separate from timeline serving.

Apache Flink processes tweet-events, like-events, and follow-events, streaming aggregated engagement metrics into ClickHouse for creator analytics dashboards and offline model feature stores.

Event Bus Design (Kafka)

The event bus decouples write ingestion from fan-out workers and background indexers.

Topic Architecture and Consumer Groups:

Topics:
  tweet-events:       New tweet creations, deletions, native retweets, and quote tweets
  like-events:        Like and unlike actions on tweets
  follow-events:      Follow and unfollow graph mutations
  mention-events:     User mention notifications and alerts

Producer Configuration:
  - tweet-events: Driven by Cassandra CDC connector tailing committed on-disk CommitLog mutations
  - like-events, follow-events, mention-events: Published by their respective owning services (Interaction Service, Social Graph Service) using reliable transactional outbox/CDC publication mechanisms

Topic Configuration (tweet-events):
  Partitions: 128
  Partition key: author_id (ensures ordered event processing per author across partition workers)
  Retention: 7 days
  Replication factor: 3, min.insync.replicas: 2

Consumer Groups:
  1. fan-out-worker:   Pushes tweet_id via ZADD to follower Redis home timelines (< 10K followers)
  2. list-fanout:      Pushes tweet_id to list timeline caches for lists containing the creator
  3. search-indexer:   Indexes tweet documents into Elasticsearch cluster
  4. notification:     Dispatches push notifications for mentions, replies, and retweets via APNs and FCM
  5. analytics:        Streams engagement events through Flink to ClickHouse

Synchronous Ingress Path:
  Validate payload -> Append tweet record to authoritative Cassandra store -> Return HTTP 201 Created

Asynchronous Event Pipeline:
  Cassandra CDC connector reads committed CommitLog mutations from disk -> Publishes event to tweet-events in Kafka -> Downstream consumer groups execute fan-out, search indexing, and push alerts without blocking the tweet creation path.
  Dead Letter Queue: tweet-events-dlq triggered after 3 retries, with alerts when consumer lag exceeds 60 seconds.

API Design

RESTful API contracts for tweet composition, chronological and ranked timeline pagination, list feeds, search, and interactions.

Post Tweet

Creates a new tweet, reply, or quote tweet with optional media attachment identifiers:

HTTP
POST /api/v1/tweets
Authorization: Bearer <user_token>
Content-Type: application/json

{
  "content": "Hello Twitter! #systemdesign",
  "media_ids": ["img-uuid-001"],
  "reply_to": null
}

HTTP/1.1 201 Created
Content-Type: application/json

{
  "tweet_id": "1541815603606036480",
  "created_at": "2026-03-13T10:00:00Z"
}

Get Home Timeline (Chronological View)

Fetches a reverse-chronological feed using 64-bit Snowflake tweet IDs as pagination cursors:

HTTP
GET /api/v1/timeline/home?type=latest&cursor=1541815603606036480&limit=20
Authorization: Bearer <user_token>

HTTP/1.1 200 OK
Content-Type: application/json

{
  "tweets": [
    {
      "tweet_id": "1541815603606036480",
      "author_id": "usr_99201481",
      "content": "Hello Twitter! #systemdesign",
      "like_count": 124,
      "retweet_count": 18,
      "is_retweet": false,
      "created_at": "2026-03-13T10:00:00Z"
    }
  ],
  "next_cursor": "1541815603606036400",
  "has_more": true
}

Get Home Timeline (Algorithmically Ranked View)

Fetches an algorithmically ranked feed using an opaque cursor encoding the score watermark, generation timestamp, and seen offsets to reduce duplicate or skipped items across pagination requests:

HTTP
GET /api/v1/timeline/home?type=ranked&cursor=eyJzY29yZSI6MTQ4LjIsImdlbl90cyI6MTcxMDMyMDAwMH0&limit=20
Authorization: Bearer <user_token>

HTTP/1.1 200 OK
Content-Type: application/json

{
  "tweets": [
    {
      "tweet_id": "1541815603606036480",
      "author_id": "usr_99201481",
      "content": "Hello Twitter! #systemdesign",
      "like_count": 124,
      "retweet_count": 18,
      "is_retweet": true,
      "retweeted_by": "usr_10293847",
      "created_at": "2026-03-13T10:00:00Z"
    }
  ],
  "next_cursor": "eyJzY29yZSI6MTI0LjYsImdlbl90cyI6MTcxMDMyMDAwMH0",
  "has_more": true
}

Get List Timeline

Retrieves tweets published by members of a specific curated list:

HTTP
GET /api/v1/lists/lst_88391028/timeline?cursor=1541815603606036480&limit=20
Authorization: Bearer <user_token>

HTTP/1.1 200 OK
Content-Type: application/json

{
  "tweets": [
    {
      "tweet_id": "1541815603606036480",
      "author_id": "usr_99201481",
      "content": "Distributed systems update #systemdesign",
      "created_at": "2026-03-13T10:00:00Z"
    }
  ],
  "next_cursor": "1541815603606036400",
  "has_more": true
}

Get User Profile Timeline

Retrieves tweets published directly by a specific user profile from Cassandra:

HTTP
GET /api/v1/users/usr_99201481/tweets?cursor=1541815603606036480&limit=20
Authorization: Bearer <user_token>

HTTP/1.1 200 OK
Content-Type: application/json

{
  "tweets": [
    {
      "tweet_id": "1541815603606036480",
      "content": "Hello Twitter! #systemdesign",
      "created_at": "2026-03-13T10:00:00Z"
    }
  ],
  "next_cursor": "1541815603606036400"
}

Search Tweets

Queries the Elasticsearch index by keyword or hashtag:

HTTP
GET /api/v1/search?q=%23systemdesign&type=recent&limit=20
Authorization: Bearer <user_token>

HTTP/1.1 200 OK
Content-Type: application/json

{
  "results": [
    {
      "tweet_id": "1541815603606036480",
      "content": "Hello Twitter! #systemdesign",
      "author_id": "usr_99201481",
      "score": 150.4
    }
  ],
  "next_cursor": "cur_993820"
}

Like and Retweet Endpoints

HTTP
POST /api/v1/tweets/{tweet_id}/like
Authorization: Bearer <user_token>

POST /api/v1/tweets/{tweet_id}/retweet
Authorization: Bearer <user_token>

Get Trending Topics

Trending computation is handled by a downstream streaming system (detailed in Trending Topics). The endpoint below illustrates the client-facing response shape:

HTTP
GET /api/v1/trends?location=US
Authorization: Bearer <user_token>

HTTP/1.1 200 OK
Content-Type: application/json

{
  "trends": [
    {
      "hashtag": "#SystemDesign",
      "tweet_count": 125000,
      "rank": 1
    }
  ]
}

Common Error Responses

400 Bad Request: invalid input, missing required fields, or malformed JSON payload
401 Unauthorized: missing or invalid authentication token or API key
403 Forbidden: authenticated caller lacks required permissions for this resource
404 Not Found: requested resource ID does not exist
409 Conflict: duplicate write or version conflict, retry with a unique idempotency key
422 Unprocessable Entity: syntactically valid request failed semantic business validation
429 Too Many Requests: rate limit quota exceeded, client should honor Retry-After header
500 Internal Error: unexpected server failure, retry safely with an idempotency key
503 Service Unavailable: downstream dependency is unavailable or overloaded, retry with exponential backoff

Data Model

Polyglot storage architecture with clearly delineated responsibilities: Cassandra serves as the authoritative time-series tweet store and author timeline repository, MySQL manages relational social graph and list metadata, Redis caches hot home timelines and tweet objects, and Elasticsearch indexes full-text search projections.

Cassandra: Authoritative Tweet Store, Lookup Index & Conversation Tables

The primary tweets table is the authoritative source of truth for tweet records, clustered chronologically per author to optimize profile timeline queries, quote tweets, and native retweets. The tweet_lookup table is a derived secondary index populated asynchronously via the CDC stream to route arbitrary GET /tweets/{tweet_id} lookups to the author's Cassandra partition. Newly created tweets requested in creator context bypass this index because user_id is already known. If an arbitrary point-lookup arrives before CDC propagation completes, the system uses a bounded read-after-write retry path rather than attempting an expensive cluster-wide Cassandra scan. The conversation_tweets table clusters reply trees by conversation_id to allow single-partition thread metadata retrieval across multiple authors:

SQL
-- Authoritative Durable Tweet Store (Cassandra)
-- Partitioned by user_id to cluster an author's tweets chronologically for profile timeline lookups
CREATE TABLE tweets (
    user_id          UUID,
    tweet_id         BIGINT,          -- 64-bit Snowflake ID (monotonically time-ordered)
    conversation_id  BIGINT,          -- Root tweet ID referencing the top of the reply thread
    content          VARCHAR(280),
    media_urls       TEXT,            -- JSON array of CDN media blob URLs
    reply_to         BIGINT,          -- Direct parent tweet_id reference (NULL if top-level tweet)
    retweet_of       BIGINT,          -- Original tweet_id if native retweet (NULL otherwise)
    quoted_tweet_id  BIGINT,          -- Quoted tweet_id if quote tweet with commentary (NULL otherwise)
    is_deleted       BOOLEAN,         -- Authoritative soft-delete watermark
    created_at       TIMESTAMP,
    PRIMARY KEY (user_id, tweet_id)
) WITH CLUSTERING ORDER BY (tweet_id DESC);

-- Derived Global Lookup Index (Cassandra / KV Store)
-- Maps tweet_id to author user_id for point lookup by tweet ID (e.g. GET /api/v1/tweets/{tweet_id})
-- Asynchronously populated via the CDC mutation stream. Lookups for newly created tweets with a known author bypass this index, while arbitrary lookups during CDC propagation use a short read-after-write retry path
CREATE TABLE tweet_lookup (
    tweet_id         BIGINT PRIMARY KEY,
    user_id          UUID
);

-- Dedicated Conversation / Thread Query Table (Cassandra)
-- Partitioned by conversation_id to enable single-partition seeks across multi-author reply trees
CREATE TABLE conversation_tweets (
    conversation_id  BIGINT,
    tweet_id         BIGINT,
    author_id        UUID,
    reply_to         BIGINT,
    created_at       TIMESTAMP,
    PRIMARY KEY (conversation_id, tweet_id)
) WITH CLUSTERING ORDER BY (tweet_id ASC);

MySQL: Social Graph Follows Table

Relational storage for creator relationships with bi-directional indexing:

SQL
-- Social Graph Storage (MySQL / Relational DB)
CREATE TABLE follows (
    follower_id  BIGINT NOT NULL,
    followee_id  BIGINT NOT NULL,
    created_at   TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP,
    PRIMARY KEY (follower_id, followee_id),
    INDEX idx_followee (followee_id, follower_id)
);

MySQL: Lists, Members, and Subscriptions Tables

Relational schema for curated list definitions, included creator rosters, and subscriber lists:

SQL
-- List Metadata, Membership, and Subscriptions (MySQL)
CREATE TABLE lists (
    list_id           BIGINT PRIMARY KEY,
    owner_id          BIGINT NOT NULL,
    name              VARCHAR(128) NOT NULL,
    description       VARCHAR(255),
    is_private        BOOLEAN DEFAULT FALSE,
    member_count      INT DEFAULT 0,
    subscriber_count  INT DEFAULT 0,
    created_at        TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP,
    INDEX idx_owner (owner_id)
);

-- Members: Creators included in the list
CREATE TABLE list_members (
    list_id           BIGINT NOT NULL,
    member_id         BIGINT NOT NULL,  -- Creator included in the list
    added_at          TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP,
    PRIMARY KEY (list_id, member_id),
    INDEX idx_member (member_id, list_id) -- Reverse index: creator -> lists containing creator
);

-- Subscribers: Users who follow/view the list feed
CREATE TABLE list_subscribers (
    list_id           BIGINT NOT NULL,
    user_id           BIGINT NOT NULL,  -- Subscriber user ID
    subscribed_at     TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP,
    PRIMARY KEY (list_id, user_id),
    INDEX idx_user_subs (user_id, list_id)
);

Redis: Home & List Timeline Sorted Set Schemas

In-memory serving cache holding up to 800 tweet IDs per active user or list:

YAML
# User Home Timeline Sorted Set (Retains up to 800 entries, subject to a sliding 48-hour TTL)
Key:     timeline:home:{user_id}
Type:    Sorted Set
Members: tweet_id
Scores:  timestamp (epoch milliseconds)
Max:     800 entries (capacity cap evicts oldest entries before sliding 48-hour TTL expires on active feeds)
TTL:     Sliding 48 hours (reset on each read)

# List Timeline Sorted Set
Key:     timeline:list:{list_id}
Type:    Sorted Set
Members: tweet_id
Scores:  timestamp (epoch milliseconds)
Max:     800 entries
TTL:     48 hours

# Celebrity Recent Tweets Serving Cache (Index for pull-on-read)
Key:     timeline:user:{celeb_id}
Type:    List / Sorted Set
Members: tweet_id (up to the latest 100 tweets subject to freshness window)
TTL:     7 days

Redis: Real-Time Interaction Counters & Trending Topics Cache

Atomic interaction counters merged on hydration, paired with trending topics populated by downstream stream analytics (see Trending Topics):

YAML
# Real-Time Interaction Counters
Key:     tweet:counters:{tweet_id}
Type:    Hash
Fields:  like_count, retweet_count, reply_count
TTL:     30 days (refreshed on interaction)

# Trending Topics Cache
Key:     trending:{country_code}
Type:    Sorted Set
Members: hashtag
Scores:  trend_velocity_score

Elasticsearch: Tweet Search Index Document

JSON
{
  "tweet_id": "1541815603606036480",
  "user_id": "usr_99201481",
  "username": "johndoe",
  "content": "Hello Twitter! #systemdesign",
  "hashtags": ["systemdesign"],
  "mentions": [],
  "created_at": "2026-03-13T10:00:00Z",
  "engagement_score": 150,
  "language": "en"
}

Fault Tolerance

Resilience strategies for fan-out worker backlogs, cache evictions, authoritative deletion synchronization, cache reconstruction, and hot-key protection.

Fault Tolerance Strategies

Failure DomainResilience Mechanism
Tweet Durability and Dual-Write LossTweets are committed to authoritative Cassandra storage before returning HTTP 201 Created. An asynchronous Cassandra CDC connector tails committed on-disk mutations to provide a durable, replayable at-least-once bridge to Kafka.
Fan-Out Worker BackpressureKafka buffers incoming tweet events across 128 partitions, allowing consumer workers to catch up gracefully after traffic surges without dropping messages.
Redis Cache Eviction & ReconstructionTimeline sorted sets and real-time interaction counters are rebuildable on demand from the follow graph, author streams, and durable interaction logs without full table scans.
Celebrity Fan-Out StormsThe hybrid architecture automatically skips push fan-out for creators exceeding 10,000 followers, eliminating linear follower-proportional write amplification.
Viral Tweet HotspotsFull tweet objects are cached in Redis with distributed read replicas and local memory caching to absorb millions of concurrent requests for viral posts.
Fan-Out Queue Lag FallbackWhen consumer queue lag spikes, the Timeline Service detects timeline staleness and temporarily falls back to fan-out on read across followees and celebrity caches to preserve freshness as much as possible until worker queues clear.

Authoritative Deletion Lifecycle

Core Invariant: The canonical tweet store in Cassandra is the sole source of truth for deletion state. Timeline caches are never trusted as the authoritative source of visibility.

  1. The Tweet Service updates the database record setting is_deleted = true as an authoritative soft delete in Cassandra.
  2. The Cassandra CDC connector captures the committed update from the on-disk CommitLog and publishes a tweet-deleted event to the Kafka event bus.
  3. The Fan-Out Service asynchronously removes the tweet_id from follower Redis home timelines using ZREM.
  4. The Search Indexer removes the corresponding document from Elasticsearch.
  5. During timeline assembly, Step 6 (Hydration) checks the is_deleted tombstone in the Redis Tweet Cache on cache hits. On cache misses or uncertainty, it falls back to the authoritative Cassandra row, ensuring deleted tweets are never rendered.
  6. For retweets of deleted tweets, hydration resolves the original tweet reference to an unavailable tombstone state.

Cache Reconstruction Strategy

If a user's Redis home timeline cache or tweet object cache is evicted or lost, the Timeline Service reconstructs it without scanning the entire historical tweet database:

  • Query the Social Graph Service for the user's current follow list.
  • Fetch recent tweet streams from those authors within a recent bounded time window (such as the last 7 days) from Cassandra.
  • Merge the tweet IDs in memory, apply initial filtering, and write up to 800 entries into timeline:home:{user_id} in Redis.
  • If Redis interaction counters (tweet:counters:{tweet_id}) are lost or evicted, counters are rebuilt on demand from durable interaction state or replayed interaction events, ensuring Redis functions as a low-latency serving accelerator rather than the authoritative source of truth.

Additional Considerations

Advanced topics covering real-time feed updates, tweet conversation trees, content moderation pipelines, and interview walkthrough strategies.

Related Problems and Concepts

The hybrid fan-out pattern is covered generically in News Feed System. Full-text search and inverted index distribution are detailed in Search Engine, while velocity-based ranking is explored in Trending Topics. Foundational scalability and cache strategies are covered in Sharding and Partitioning, Caching Patterns and Invalidation, Scaling 0 to 1M Users, and System Design Interview Patterns.

Real-Time Feed Updates

Rather than streaming every newly published tweet over active WebSockets and saturating client bandwidth, the client maintains a Server-Sent Events (SSE) connection. The server pushes lightweight notification badges (such as "12 new tweets"), prompting the user to pull and refresh their feed on demand.

Tweet Thread and Conversation View

Replies form a recursive tree structure linked via reply_to tweet-ID references rather than relational foreign keys. Because the primary tweets table is partitioned by (user_id, tweet_id), querying a conversation across multiple authors requires a dedicated query model. The system maintains a conversation_tweets table in Cassandra partitioned by conversation_id (the root tweet ID) with tweet_id as clustering key. Thread retrieval fetches candidate reply IDs and metadata from this table, followed by hydration against the authoritative Cassandra tweets table (or warm Tweet Cache). Validating the authoritative is_deleted flag during hydration ensures deleted replies are rendered as tombstones or omitted from the tree, preventing phantom entries from lingering conversation index rows.

Content Moderation

Pre-publish classifiers evaluate text and media attachments for spam and policy violations. Post-publish signals, including user reports and engagement velocity spikes, route flagged content into automated shadow-banning filters or human moderation review queues.

Interview Walkthrough

  • 25-Minute Interview Strategy

    Prioritize the hybrid fan-out decision, Cassandra CDC streaming, and 7-stage timeline merge pipeline before discussing search indexing and analytics.

    • Fan-out decision: push vs pull vs hybrid trade-offs (5 min)
    • Cassandra CommitLog CDC write path and Kafka streaming (5 min)
    • Celebrity threshold, list feeds, and Redis cache design (5 min)
    • 7-stage timeline merge, ranking, and deduplication (6 min)
    • Capacity math and write amplification calculations (4 min)
  • Open with the fan-out on write vs fan-out on read dilemma, explaining how this single choice dictates storage, latency, and celebrity handling across the architecture.
  • Propose the hybrid model early: precompute timelines for normal creators (<10K followers) while pulling celebrity tweets at read time.
  • Detail the Cassandra CommitLog CDC pattern to demonstrate senior awareness of dual-write failure modes between the database and message broker.
  • Walk through the 7-stage timeline merge step-by-step: Candidate Generation (~300 candidates), Visibility Filtering, Celebrity Merge, Retweet Deduplication, Ranking, Hydration (top 20), and Response.
  • Explain native retweet storage (referencing original tweet ID) and deduplication during feed assembly.
  • Highlight write amplification math (500M tweets x 200 avg followers = 100B fan-out writes daily) to justify the adaptive 10,000 follower threshold.

Engineering Trade-offs

Key architectural trade-offs across timeline assembly, celebrity threshold calibration, event publication consistency, and pagination mechanics.

Home Timeline Assembly: The 7-Stage Read Walkthrough

Scenario: User B opens the home feed. User B follows 300 normal creators and 5 celebrities.

Stage 1: Candidate Generation (Illustrative ~1 to 2 ms budget)
  - Precomputed Normal Tweets: ZREVRANGEBYSCORE timeline:home:{B} +inf -inf LIMIT 0 200
    Returns ~200 recent tweet IDs from Redis home timeline cache (which stores up to 800 entries subject to a sliding 48-hour TTL).
  - Celebrity Pull: 5 parallel Redis queries (ZREVRANGEBYSCORE timeline:user:{celeb_id} {now} {now - 24h} LIMIT 0 20)
    Returns ~100 recent tweet IDs from the 5 followed celebrities via fan-out-on-read (up to latest 100 tweets per celebrity within the freshness window).
  - Total working candidate pool: ~300 candidate tweet IDs.

Stage 2: Visibility Filtering (Illustrative < 1 ms budget)
  - Checks User B's cached follow graph, block set, and mute lists in Redis.
  - Drops candidates from unfollowed, blocked, or muted accounts.

Stage 3: Celebrity Merge (Illustrative ~2 ms budget)
  - Merges 200 precomputed candidates + 100 celebrity candidates (300 total) using an in-memory multi-way merge with O(C log S) complexity where C is candidates consumed and S is sorted streams.

Stage 4: Retweet Deduplication (Illustrative < 1 ms budget)
  - Collapses duplicate original/retweet candidate pairs into a single entry with a retweet attribution badge.

Stage 5: Ranking & Scoring (Illustrative ~5 ms budget)
  - For algorithmic feeds: Gradient-boosted decision tree scores each of the remaining ~280 candidates:
    score = (alpha * recency) + (beta * engagement) + (gamma * affinity) + (delta * content_type)
    Sorts candidates and selects the top 20 items.
  - For chronological feeds: Skips ML scoring and sorts directly by timestamp ("Latest Tweets").

Stage 6: Hydration (Illustrative ~3 ms budget)
  - Fetches full tweet objects for ONLY the top 20 candidates from Tweet Object Cache (Redis Hash: HGETALL tweet:{tweet_id}) and real-time counts from Redis Hash (tweet:counters:{tweet_id}).
  - Validates is_deleted = false from cache tombstone; falls back to authoritative Cassandra row on cache miss or invalidation.
  - If deleted, drops the item and takes the next candidate.

Stage 7: Response Assembly (Illustrative < 1 ms budget)
  - Encodes pagination cursor (opaque ranking watermark or chronological tweet_id) and returns the top 20 hydrated tweets.
  - Total illustrative internal serving budget: ~12 to 15 ms under healthy warm-cache conditions with local network hops (distinct from the formal production SLO of p99 < 200 ms, which accommodates cold caches, cross-region replication, and peak load).

The Celebrity Threshold: Dynamic Workload Calibration

Fan-Out Write Cost per Tweet (Illustrative Sizing Assumptions with Pipelining):
  Write operations per tweet = 1 tweet * N followers (Redis ZADD operations)

  - 100 followers:    100 ZADDs -> < 1 ms (illustrative sizing assumption)
  - 10,000 followers: 10,000 ZADDs -> ~10 ms (illustrative sizing assumption with pipelining)
  - 100,000 followers: 100,000 ZADDs -> ~100 ms (begins creating consumer queue lag)
  - 50,000,000 followers: 50,000,000 ZADDs -> ~50 seconds (unacceptable linear write amplification)

Threshold Evaluation:
  - Below 10K followers: Accounts for 99.9% of users, and fan-out completes within acceptable worker queue latency.
  - Above 10K followers: Accounts for ~0.1% of users, but generates over 80% of potential write amplification.

Adaptive Threshold Control:
  - Off-Peak Traffic (e.g. overnight): Threshold can be raised to 50K followers to maximize precomputation.
  - Peak Live Events (e.g. Super Bowl, breaking news): Threshold dynamically lowers to 5K followers to protect fan-out worker queues.

Queue Lag Fallback:
  If fan-out consumer lag grows during unexpected traffic surges, the Timeline Service detects timeline staleness and automatically falls back to fan-out on read across followees and celebrity caches to preserve freshness as much as possible until queues clear.

Dual-Write Mitigation: Cassandra CommitLog CDC vs Direct Dual-Write

The Dual-Write Problem:
  Tweet Service writes to Database -> Database write succeeds.
  Tweet Service publishes to Kafka -> Kafka network timeout or broker failure occurs.
  Result: Tweet exists in database but is never fanned out to follower timelines (Ghost Tweet).

Cassandra CommitLog CDC Solution:
  1. Tweet Service writes tweet record directly to authoritative Cassandra table.
  2. Cassandra appends mutation to its sequential on-disk CommitLog and returns write success.
  3. Tweet Service returns HTTP 201 Created immediately to the user.
  4. Cassandra CDC connector tails the commit log on disk and publishes events to Kafka with connector-level at-least-once delivery (subject to connector health and commit-log retention).
  5. Downstream workers process events idempotently using natural ZADD idempotency and tweet_id deduplication.

Feed Pagination: Chronological Snowflake vs Opaque Ranking Cursors

Chronological Feed Pagination (Latest Tweets):
  - Uses the 64-bit Snowflake tweet_id of the last item on the page as cursor:
    GET /api/v1/timeline/home?type=latest&cursor=1541815603606036400&limit=20
  - For the home timeline, the next page resumes the Redis home-timeline sorted set below the cursor's timestamp/Snowflake boundary (ZREVRANGEBYSCORE timeline:home:{user_id} (cursor_score -inf LIMIT 0 20), merging with celebrity cache streams).
  - For direct user profile timelines, the next page queries Cassandra partition rows where tweet_id < cursor.
  - Deterministic, stateless, and naturally immune to drift.

Algorithmically Ranked Feed Pagination (Home Feed):
  - Because ML ranking scores change dynamically as engagement accrues, simple offset or tweet ID pagination causes duplicate tweets or skipped content across pages.
  - Solution: Opaque cursor encoding ranking score watermark, generation timestamp, and seen item offsets:
    GET /api/v1/timeline/home?type=ranked&cursor=eyJzY29yZSI6MTQ4LjIsImdlbl90cyI6MTcxMDMyMDAwMH0&limit=20
  - Carries ranking watermark and state to reduce duplicate or skipped items across pages without requiring heavy server-side session state (best-effort consistent under dynamic feature changes, while chronological pagination is naturally deterministic).

💬Review

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...