Interview Setup
Interview Prompt
Design a distributed social graph store that supports follow and unfollow mutations, follower listings, mutual friend queries, and two-hop friend-of-friend lookups across 2B users and 500B edges.
Clarifying Questions (ask before designing)
| Question | Why it matters |
|---|---|
| What query patterns matter: 1-hop only or multi-hop traversals? | Over 99% of production queries are single-hop lookups (retrieving followers or following). Multi-hop exploration operates as an offline batch process at this scale. |
| Are edges directed or bidirectional? | Follow actions are directed, whereas friendships are bidirectional (stored as two directed edges or one undirected row), directly impacting write amplification. |
| Should we choose a native graph database (Neo4j) or an adjacency list (Cassandra)? | Graph databases excel at multi-hop traversals but struggle past 100B edges. The Cassandra and TAO pattern has proven scalability at multi-billion-edge production workloads. |
| How do celebrity accounts with 100M+ followers affect partition distribution? | A single followee partition becomes an extreme hot spot, making bucket sharding by follower ID prefix mandatory. |
Scope
In scope
- Adjacency list storage in Cassandra (TAO pattern)
- Single-hop queries (followers, following, is-following checks)
- Two-hop offline graph computation (friends-of-friends via Spark)
- Graph partitioning and bucket sharding for celebrities
- Mutual friends computation
- Bidirectional friendship edge semantics
Out of scope (state explicitly)
- Follow and unfollow product API and feed fan-out
- Full news feed fan-out
- Graph machine learning recommendation models
- Real-time global influence scoring
Functional Requirements
Clarify directed vs undirected edge semantics and query patterns with your interviewer. The design must accommodate follow and unfollow actions, follower and following listings, and mutual friend lookups, while confirming the scope of multi-hop graph traversals. Storage model decisions directly dictate whether queries can complete within sub-50ms latency targets.
In the room: two-hop "friends of friends" traversals are not viable as real-time online queries across billions of edges, so state clearly that you precompute traversals with Spark and cache hot results.
- Follow and Unfollow: User A creates or deletes a directed edge to User B.
- Bidirectional Friendship: Support mutual follow confirmations and symmetric connections.
- Follower List Retrieval: Fetch paginated lists of users who follow a given account.
- Following List Retrieval: Fetch paginated lists of accounts that a given user follows.
- Mutual Friends Calculation: Surface shared connections (such as "You and Alice have 12 mutual friends") in real time.
- Friend-of-Friend Suggestions: Discover second-degree connections for social recommendations.
- Graph Analytics Queries: Support shortest-path analysis, connected component discovery, and influence scoring offline.
- Account Blocking: Strictly exclude blocked users from relationship listings and suggestion algorithms.
Non-Functional Requirements
The system must deliver follow and unfollow operations in under 100 ms and provide paginated list reads at massive scale, with dedicated mitigations for celebrity hot keys.
- Low Latency: Return follower and following lists in under 50 ms p99.
- High Write Throughput: Sustain hundreds of millions of follow and unfollow operations daily.
- Scale: Scale to 2B+ registered users and 500B+ relationship edges.
- Immediate Read-Your-Writes Consistency: A follow action must be immediately reflected to the acting user.
- High Availability: Guarantee 99.99% uptime for core social graph queries.
- Traversal Performance: Complete cached two-hop queries in under 200 ms.
Capacity Estimations
Billion-edge graphs require sharding, so run this math before picking Cassandra vs a graph database. Edge count multiplied by bytes per edge tells you partition count, while celebrity fan-out drives hot-key mitigation rather than just raw storage capacity.
| Metric | Calculation | Value |
|---|---|---|
| Users (nodes) | Given | 2B |
| Edges (follow relationships) | Given | 500B |
| Follow/unfollow ops / day | Given | 500M |
| Follower list queries / sec | Derived from daily volume ÷ 86400 (+ peak factor) | 100K |
| Edge record size | Given | 32 bytes |
| Total edge storage | Given | 16 TB |
Architecture Diagram
The architecture shards the graph by user ID, maintains adjacency lists across dual Cassandra tables (following and followers), and indexes reverse edges so that checking who follows a user and who a user follows are both O(1) partition lookups. Celebrity accounts use bucket-sharded follower lists, while multi-hop traversals run offline in Apache Spark.
In the room: distinguish this design from the Follower and Following System, which handles product API endpoints, because here you are designing the underlying distributed graph storage layer.
Component Deep Dives
Storage: Adjacency List vs Edge Table
Storage model choices dictate throughput and query boundaries in a social graph. The trade-offs below examine edge tables, Cassandra dual adjacency tables, and native graph databases.
Approach 1: Edge Table (MySQL) CREATE TABLE follows (follower_id BIGINT, followee_id BIGINT, PRIMARY KEY (follower_id, followee_id)); ✓ Simple, strong consistency ✗ 2-hop traversals require JOIN -> expensive at scale Approach 2: Adjacency List in Cassandra (Facebook TAO-inspired) Two tables: following(user_id, follows) + followers(user_id, follower) ✓ Scales horizontally, fast single-hop queries ✗ Dual writes needed Approach 3: Graph DB (Neo4j, Neptune) ✓ Native graph traversals in < 50ms ✗ Harder to scale beyond 100B edges
Facebook TAO: The Industry Standard
Meta designed TAO to scale association reads and writes across billions of edges by layering a read-through distributed cache over MySQL.
TAO (The Associations and Objects): Facebook's custom graph store Core abstraction: Objects: users, pages, groups, posts (nodes) Associations: follows, likes, friendships (edges) Stored in MySQL, cached aggressively in distributed cache layer: Cache hit rate: > 99.8% Read from cache: < 1ms Write path: 1. Write to MySQL (source of truth) 2. Invalidate cache for both id1 and id2 3. Async: update denormalized count tables Key insight: 99%+ of queries are single-hop (get followers/get following). Only friend suggestions need multi-hop -> done offline in Spark.
Event Bus (Kafka)
Graph mutations invalidate the TAO cache synchronously and publish events to graph-events so that notifications, downstream analytics, and offline suggestion pipelines never block the critical follow path.
Topic: graph-events
Partitions: 64
Partition key: source_user_id (routes edge-change events consistently)
Retention: 7 days
Replication factor: 3, min.insync.replicas: 2
Producer: Social Graph Service after primary store write + cache invalidation
Event: { event_id, action: "follow" | "unfollow" | "friend", source_id, target_id, timestamp }
Consumer groups:
1. notification: asynchronous push and email notifications on new edges
2. count-reconciler: periodic counter validation against raw edge tables
3. suggestion-batch: incremental mutual-friend index updates
4. analytics: graph growth metrics streamed to ClickHouse
Sync path: write edge, invalidate TAO cache, publish graph-events, and return 200 OK.
Async path: notification delivery and count reconciliation proceed without blocking graph writes.
DLQ: graph-events-dlq, alerting when consumer lag exceeds 60 seconds.Mutual Friends: The Interview Favorite
Computing mutual connections between two users is an interview staple requiring fast set operations in memory.
"Alice and Bob have 12 mutual friends": how to compute? Optimization: Redis sorted set intersection ZINTERSTORE mutual_temp followers:alice followers:bob ZCARD mutual_temp -> count Redis does this in-memory in < 5ms for sets up to 10K members. For celebrities (100M followers): Don't compute real-time. Pre-compute mutual count in batch. Or show: "You follow [3 friends who follow this celebrity]"
API Design
Social Graph Endpoints
The API provides endpoints for follow and unfollow operations, cursor-paginated follower and following lists, and mutual friend calculations.
POST /api/v1/follow
{ "target_user_id": "bob" } -> 200 OK
DELETE /api/v1/follow
{ "target_user_id": "bob" } -> 200 OK
GET /api/v1/users/{uid}/followers?cursor=...&limit=50
-> { "followers": [{id, name, avatar}, ...], "cursor": "..." }
GET /api/v1/users/{uid}/following?cursor=...&limit=50
-> { "following": [...], "cursor": "..." }
GET /api/v1/users/{uid}/mutual-friends?with=bob
-> { "mutual": [{id, name}, ...], "count": 12 }Common Error Responses
Structured error schemas handle blocked relations, private account restrictions, and invalid target identities.
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
Relational and Cassandra Schema: Source of Truth
Relational stores maintain primary edge entries, while Cassandra stores denormalized following and followers tables to optimize single-hop query paths. Review Sharding and Partitioning for distribution details.
CREATE TABLE follows (
follower_id BIGINT, followee_id BIGINT,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (follower_id, followee_id)
);
CREATE INDEX idx_followee ON follows (followee_id, follower_id);Redis Cache Layer
Redis Sorted Sets maintain paginated follower and following lists ordered by follow timestamp, with companion hash structures tracking aggregate relationship counters.
followers:{uid} -> Sorted Set (score=timestamp, member=user_id)
following:{uid} -> Sorted Set (score=timestamp, member=user_id)
follow_count:{uid} -> Hash {followers: 1523, following: 342}Fault Tolerance
| Concern | Solution |
|---|---|
| Dual write consistency | Write both tables within the same Cassandra batch or MySQL transaction to prevent asymmetric edges. |
| Counter drift | Run an asynchronous hourly reconciliation job comparing counter values against raw edge tables. |
| Cache invalidation | On follow or unfollow mutations, invalidate the cached sorted sets for both involved user IDs. |
| Celebrity hot partition | Shard followers across compound keys using (followee_id, follower_id_prefix) bucketing. |
Race Conditions: Follow and Unfollow in Quick Succession
When users rapidly toggle follow state, network reordering can cause stale overwrites if timestamps are omitted.
T=0: Follow Bob -> INSERT follows (alice, bob) T=50ms: Unfollow Bob -> DELETE follows (alice, bob) If out of order at DB: DELETE arrives first (no-op, row doesn't exist) INSERT arrives second -> alice follows bob (WRONG!) Solution: Include timestamp, use LWW (Last-Writer-Wins). Or use Cassandra (naturally LWW with cell-level timestamps).
Additional Considerations
Related Problems
The Follower and Following System covers the product API layer: idempotent follow and unfollow operations, is-following checks at 500K QPS, and feed fan-out thresholds. This article focuses on the graph storage layer (the TAO pattern, 500B edges, celebrity bucket sharding, and offline two-hop traversals). You can discuss them as separate modules or stack the product API on top of this graph storage layer.
Interview Walkthrough
- 25-minute pacing strategy
Prioritize core dual-table sharding and cache invalidation before discussing offline recommendation pipelines.
- Draw dual Cassandra tables (following and followers) and explain why both directions are hot paths (5 min)
- Explain single-hop follow queries are O(1) while two-hop traversals need offline Spark precomputation (6 min)
- Quantify 500B edges and ~16 TB storage: shard by user_id with celebrity bucket fan-out (5 min)
- Cover graph databases vs adjacency lists: evaluate when Neo4j operational cost is justified (5 min)
- Staff depth: blocking edges, influence scoring, and cache warming for celebrity follower lists (4 min)
- Model the social graph as directed edges from follower to followee stored in Cassandra with indexes on both directions.
- Apply the TAO pattern: Cassandra or MySQL as source of truth, with Redis as a cache-aside layer for hot single-hop queries under 1 ms.
- Cache follower and following lists in Redis sorted sets paginated by cursor, invalidating the cache entry on follow and unfollow writes.
- Maintain denormalized edge counts in Redis hashes rather than running expensive COUNT(*) queries across billions of rows on the read path.
- For mutual-friend queries, use bidirectional BFS or pre-computed second-degree indexes, because unidirectional BFS at depth 3 visits millions of nodes.
- Handle concurrent follow and unfollow mutations with Last-Writer-Wins timestamps to prevent out-of-order writes from leaving stale edges.
- Explain why graph databases such as Neo4j struggle at 500B+ edges: Facebook built TAO specifically because off-the-shelf graph databases could not scale to their operational throughput.
- Common pitfall: proposing Neo4j or a graph database for a billion-user social network, when interviewers expect the TAO relational storage and distributed cache pattern instead.
Engineering Trade-offs
Architectural Trade-offs
Social graph architectures balance adjacency list partitioning against edge stores, pull vs push fan-out mechanics, and cache invalidation consistency.
Graph DB vs Relational + Cache
For social networks at scale: Relational + cache (TAO pattern). For knowledge graphs / small-scale: Graph DB. At Facebook scale: 2B users x 250 avg connections = 500B edges. No graph DB scales to this level. TAO: MySQL + massive cache layer -> custom-built for their scale.
Degree of Separation: Bidirectional BFS
Why bidirectional?
Unidirectional BFS at depth 3: visits ~200^3 = 8M nodes
Bidirectional BFS at depth 3: visits ~2 x 200^1.5 = ~5.6K nodes
-> 1400x fewer nodes explored!
LinkedIn's approach:
Pre-compute 1st and 2nd degree connections offline (Spark)
Store in graph index: 2nd_degree:{uid} = [user_ids]Graph Partitioning: Sharding Edges
Option 1: Hash partition by user_id ✓ Even distribution ✗ Cross-shard queries for friend-of-friend Option 2: Social-aware partitioning (METIS) ✓ Most traversals stay within one shard ✗ Complex to compute, must re-partition periodically Option 3: Hash + replicated edge list + aggressive caching Facebook TAO: hash partition by user_id with 99.8% cache hit rate
Review
How helpful was this walkthrough?
Click a star to rate. We actively use this feedback to refine and update our system design content.
Discussion
Share your thoughts, ask questions, or help others.