A social graph is the data structure behind "friends", "followers" and "connections" on a social product. Each user is a node (also called a vertex), and each relationship is an edge between two nodes. The graph answers questions like "who are my friends?", "how am I connected to this person?" and "who might I know?".
This is a favourite interview problem because it looks like a textbook graph exercise but breaks the moment the graph stops fitting on one machine. A breadth-first search that takes microseconds in memory becomes thousands of network calls when friends live on different servers. Interviewers use it to probe whether you can reason about fan-out, sharding, caching and the gap between an algorithm and a distributed system.
What interviewers typically probe:
- How you store edges (adjacency lists in a key-value or relational store, or a graph database) and why.
- How you find the shortest path between two users when the graph is spread across many shards.
- How "people you may know" works without computing it on every page load.
- How you handle celebrity accounts with millions of followers.
If you need a refresher on sharding and replication first, read Databases at scale and Caching.
Problem and scope
You are designing the relationship layer of a social network. It is a backend service that other features (feed, search, messaging, profile pages) call. It does not render pages or store posts.
In scope:
- Create and remove relationships between users.
- List a user's friends or followers, paginated.
- Find the shortest path (degrees of separation) between two users, for example "You → Asha → Ravi → Meera".
- Suggest "people you may know" (PYMK).
Out of scope: news feed ranking, messaging, privacy settings beyond simple visibility, and the ads graph. Say these out loud in an interview so the interviewer can redirect you if they want one of them.
Clarifying questions
Asking good questions narrows the design. Here are the ones that change the answer, with the assumptions this lesson uses.
| Question | Why it matters | Assumption here |
|---|---|---|
| Are relationships mutual (friend) or one-way (follow)? | Mutual edges must be written twice and kept in sync; follows are directional | Support both; friendship is mutual, follow is one-way |
| How many users and how many connections each? | Drives storage and BFS fan-out | 500 million users, average 200 connections |
| Is there a cap on connections? | A cap bounds fan-out | Friends capped at 5,000; followers uncapped |
| How fresh must "degrees of separation" be? | Exact live answer vs precomputed | Live, but may be a few seconds stale |
| What is the maximum path length we search? | Bounds the cost of a search | Stop at 6 hops and report "not closely connected" |
| Is PYMK real-time? | Online computation vs batch | Batch daily, plus a light real-time boost |
Interview tip
Ask about mutual vs one-way edges first. It changes the data model, write path and consistency story. A follow graph (one-way) is asymmetric: a celebrity may have 50 million followers but follow 200 people, and that asymmetry drives most of the hard design choices.
Functional and non-functional requirements
Functional requirements describe what the system does:
addFriend(a, b)after a request is accepted;removeFriend(a, b).follow(a, b)andunfollow(a, b).listFriends(user, cursor)andlistFollowers(user, cursor), sorted by most recent first.areConnected(a, b)returns a yes/no in constant time.shortestPath(a, b, maxHops)returns one shortest path or "none within maxHops".suggestions(user)returns up to 50 ranked PYMK candidates with a reason ("12 mutual friends").
Non-functional requirements describe how well it does it:
- Low latency for the hot reads:
areConnectedandlistFriendsunder 20 ms at the 99th percentile (p99, meaning 99 out of 100 requests are faster). - Shortest path under about 1 second at p99; it is a less frequent, user-initiated feature.
- High availability for reads: profile pages and feeds call this service constantly. A brief outage would break much of the product.
- Durability: a lost friendship is a visible bug and a support ticket.
- Eventual consistency is acceptable between the two directions of a friendship for a few seconds, but they must converge.
Back-of-the-envelope estimates
These numbers decide whether the graph fits on one machine (it does not) and how expensive a search is.
Edges. 500 million users × 200 connections = 100 billion directed edge rows. A friendship between A and B is stored as two rows (A→B and B→A) so that "list A's friends" and "list B's friends" are both single lookups. That means about 50 billion friendships.
Storage. One edge row needs at least the two user IDs. With 64-bit (8-byte) IDs that is 16 bytes, so:
- Raw: 100 × 10⁹ × 16 bytes = 1.6 TB.
- Add a created-at timestamp, row overhead and an index, and assume roughly double: about 3.2 TB.
- Three replicas for durability and read capacity: about 9.6 TB.
That is too large for one server's memory and comfortable across a few dozen shards (a shard is one partition of the data living on its own server group).
Per user. One adjacency list of 200 IDs is 200 × 8 = 1.6 KB. Reading a whole friend list is a single small read. That fact makes adjacency lists the natural storage unit.
Write rate. Assume 100 million new friendships or follows per day. 100 × 10⁶ ÷ 86,400 ≈ 1,157 per second, and because each mutual friendship writes two rows, about 2,300 row writes per second on average. Peaks of 3× give roughly 7,000 row writes per second, spread across shards. Modest.
Read rate. Assume 40% of users are daily active and each triggers 20 friend-list or areConnected reads (profile visits, feed filtering, privacy checks). 500M × 0.4 × 20 = 4 billion reads per day ≈ 46,000 reads per second on average, well over 100,000 at peak. Reads dominate, so caching matters.
Shortest path queries. Assume 1 billion per month. 10⁹ ÷ (30 × 86,400) ≈ 386 per second, about 1,160 per second at a 3× peak. Each one is expensive, so the per-query cost is what matters.
The fan-out problem. With an average degree (number of connections) of 200, a plain BFS from one user touches:
| Hops from the source | Users reached (upper bound, 200ᵈ) |
|---|---|
| 1 | 200 |
| 2 | 40,000 |
| 3 | 8,000,000 |
| 4 | 1,600,000,000 |
Real graphs overlap heavily (your friends know each other), so actual counts are lower, but the shape is right: every hop multiplies work by roughly the degree. By three hops a one-sided search is touching millions of users, most of whom live on other shards.
Bidirectional search to the rescue. If you search from both ends and meet in the middle, a 4-hop path needs each side to go only 2 hops: 2 × 200² = 80,000 users instead of 1.6 billion. For a 6-hop path it is 2 × 200³ = 16 million versus 200⁶ = 64 trillion. This single idea is the core of the shortest-path design.
Common mistake
Do not quote "six degrees of separation" as a reason the search is cheap. Short paths exist, but finding one by brute force is still exponential in path length. The answer is to cut the exponent in half with bidirectional search and to cap the hop count.
API design
An internal service is usually called over gRPC or HTTP with JSON. The shapes matter more than the transport. See API design for pagination and idempotency patterns.
POST /v1/friendships body: {"user_a": 42, "user_b": 77}
DELETE /v1/friendships/42/77
POST /v1/follows body: {"follower": 42, "followee": 9}
DELETE /v1/follows/42/9
GET /v1/users/42/friends?limit=50&cursor=<opaque>
-> {"friends": [77, 103, ...], "next_cursor": "..."}
GET /v1/users/42/connected/77
-> {"connected": true, "since": "2026-03-01T10:00:00Z"}
GET /v1/path?from=42&to=901&max_hops=6
-> {"path": [42, 77, 310, 901], "hops": 3}
-> {"path": null, "reason": "not_within_max_hops"}
GET /v1/users/42/suggestions?limit=20
-> {"suggestions": [{"user": 555, "mutual": 12}, ...]}
Design notes:
- Cursor pagination (an opaque token encoding the last seen
(created_at, friend_id)) rather than page numbers. Offsets get slower as you page deeper and skip or repeat rows when friends are added mid-scroll. - Friendship writes are idempotent: adding an existing friendship returns success without creating duplicates. Clients retry on timeouts, and retries must be safe.
- The path API takes
max_hopsso callers cannot request an unbounded search.
Data model and storage choice
Adjacency lists in a key-value or relational store
An adjacency list stores, for each user, the list of users they connect to. In a relational database that becomes one row per directed edge:
CREATE TABLE edges (
src_id BIGINT NOT NULL, -- owner of the list
edge_type SMALLINT NOT NULL, -- 1 = friend, 2 = follows, 3 = followed_by
dst_id BIGINT NOT NULL,
created_at TIMESTAMP NOT NULL,
PRIMARY KEY (src_id, edge_type, dst_id)
);
-- Serves "list friends newest first" without a sort.
CREATE INDEX edges_by_time ON edges (src_id, edge_type, created_at DESC);
Key ideas:
- Primary key
(src_id, edge_type, dst_id)makesareConnected(a, b)a single point lookup and prevents duplicate edges. - Every row starts with
src_id, so if you shard bysrc_id, a user's entire list lives on one shard. Listing friends never crosses shards. - One-way follows are stored twice too:
A follows Bon A's shard andB followed_by Aon B's shard. Otherwise "list B's followers" would have to scan every shard.
In a wide-column store such as Cassandra or HBase, the same idea is a partition per src_id with clustering by (edge_type, created_at, dst_id). In a plain key-value store you can keep friends:<user_id> as a list or sorted set, though very large lists then need chunking.
Graph databases
A graph database (Neo4j, Amazon Neptune, JanusGraph and others) stores nodes and edges natively and offers traversal queries such as "find shortest path". They shine when queries walk many hops with varied edge types, such as fraud rings or knowledge graphs.
| Aspect | Adjacency lists in KV / SQL | Graph database |
|---|---|---|
Point reads (areConnected, list friends) | Excellent, one key lookup | Good |
| Multi-hop traversal | Application code issues one batch per hop | Built-in query language |
| Horizontal sharding at 100B edges | Mature, well understood | Harder; cross-partition traversals are costly |
| Operational familiarity | High | Lower in most teams |
| Best fit | Huge, simple, read-heavy social graphs | Rich, smaller graphs with deep queries |
Large social networks have generally built graph services on top of sharded key-value or relational storage with a caching layer, rather than on a single off-the-shelf graph database. The reason is simple: 95% of traffic is one-hop reads, which adjacency lists serve best, and the rare multi-hop query can be written as a few batched rounds in application code.
Interview tip
Say: "Most traffic is one hop, so I optimise storage for one-hop reads with adjacency lists sharded by user. Multi-hop queries are rare and I implement them as batched BFS rounds in a dedicated service." That shows you chose by access pattern, not by the word "graph".
Choosing a shard key
Shard by src_id using a hash of the user ID (or consistent hashing, so adding shards moves only a fraction of keys). Every friend list is local to one shard, and lookups go straight to the right shard. The cost: a friendship write touches two shards (A's and B's), and a BFS hop touches many shards. We address both below.
High-level design
clients / other services
|
+------+------+
| API gateway |
+------+------+
|
+-------------+--------------+
| | |
+----v-----+ +-----v------+ +-----v-------+
| Graph | | Path | | Suggestion |
| service | | service | | service |
+----+-----+ +-----+------+ +-----+-------+
| | |
| batched multi-get | reads precomputed
v v v
+----------------------+ +--------------+
| Edge cache (Redis / | | PYMK store |
| Memcached), by user | | (KV, by user)|
+----------+-----------+ +------^-------+
| miss |
+----------v-----------+ |
| Edge shards 1..N | |
| (SQL or wide-column) | |
+----------+-----------+ |
| change events (CDC) |
+----------v-----------+ +------+-------+
| Event log (Kafka) +-->| Batch / ML |
+----------------------+ | PYMK jobs |
+--------------+
Components:
- Graph service: owns writes and one-hop reads. Knows the shard map.
- Edge cache: holds hot adjacency lists in memory, keyed by user ID.
- Edge shards: the source of truth, each a replicated primary with followers.
- Path service: runs bidirectional BFS by calling the graph service in batches.
- Event log: every edge change is published (via change data capture, CDC, which streams committed database changes as events) for caches, PYMK and analytics. See Transactions and CDC.
- Suggestion service: serves PYMK lists computed offline.
Request flows
Flow 1: accept a friend request
- Client calls
POST /v1/friendshipswith A and B. - Graph service validates (both users exist, not blocked, under the friend cap).
- It writes
A→Bon A's shard andB→Aon B's shard. Because these are different shards, there is no single local transaction. The service writes an intent record first (an outbox row on A's shard: "create B→A"), then a worker applies the second half and retries until it succeeds. Both inserts are idempotent thanks to the primary key. - Each committed write emits a change event.
- Cache invalidation consumers delete
friends:Aandfriends:Bfrom the cache (delete, not update, so the next read reloads the truth). - PYMK pipelines consume the event; the real-time booster removes B from A's suggestions immediately.
For a few hundred milliseconds A may see B as a friend while B does not yet see A. That brief asymmetry is acceptable under our requirements, and a nightly reconciler scans for one-sided edges and repairs them.
Flow 2: list friends
GET /v1/users/42/friends?limit=50.- Graph service looks up
friends:42in cache. On a hit, slice and return. - On a miss, read from the shard that owns user 42 using the
(src_id, edge_type, created_at DESC)index, fill the cache, and return.
Flow 3: degrees of separation
GET /v1/path?from=A&to=J&max_hops=6.- Path service runs bidirectional BFS (next section). Each BFS level becomes one batched call: "give me the friend lists of these 300 users".
- The graph service groups those IDs by shard, sends one multi-get per shard in parallel, merges results and returns them.
- When the two frontiers meet, the path service rebuilds the path from parent pointers and returns it.
Deep dive 1: shortest path with bidirectional BFS
Plain BFS recap
Breadth-first search (BFS) explores a graph level by level: first all nodes 1 hop away, then 2 hops, and so on. In an unweighted graph the first time BFS reaches the target, it has found a shortest path. Each node remembers its parent (the node it was discovered from) so you can walk back to the source. The cost is proportional to everything within the target's distance, which, as the estimates showed, explodes with hop count.
Bidirectional BFS
Run two BFS searches at once: one forward from the source, one backward from the target. Always expand the smaller frontier (the set of nodes discovered in the latest level), because it costs fewer lookups. Stop as soon as a node discovered by one side has already been seen by the other. Join the two parent chains at that meeting node.
Why it works: a path of length L is found when the forward side has gone about L/2 hops and the backward side about L/2 hops. Work drops from roughly bᴸ to 2 × b^(L/2), where b is the degree.
Worked example
Take this small graph of ten users. Lines are friendships.
A ----- B
| |
C ----- D ----- F
| |
E ----- G ----- H
| |
J ----- I
Edges: A-B A-C B-D C-D C-E D-F
E-G F-H G-H H-I G-J I-J
Adjacency lists (what one shard lookup returns):
| User | Friends |
|---|---|
| A | B, C |
| B | A, D |
| C | A, D, E |
| D | B, C, F |
| E | C, G |
| F | D, H |
| G | E, H, J |
| H | F, G, I |
| I | H, J |
| J | G, I |
Goal: shortest path from A to J.
Round 1. Forward frontier and backward frontier are both size 1, so expand forward. Fetch A's list: B, C. Neither has been seen by the backward side (which has only seen J). Forward visited: A, B, C, with parents B←A and C←A. Forward frontier is now .
Round 2. Forward frontier has 2 nodes, backward has 1, so expand the smaller one: backward. Fetch J's list: G, I. Neither is in the forward visited set. Backward visited: J, G, I, with parents G←J and I←J. Backward frontier is .
Round 3. Both frontiers have 2 nodes; ties go forward. Fetch lists for B and C in one batch.
- B: A (already visited), D (new). D is not in backward visited. Parent D←B.
- C: A (visited), D (visited), E (new). E is not in backward visited. Parent E←C.
Forward frontier is .
Round 4. Both size 2; expand forward. Fetch D and E.
- D: B, C (visited), F (new, not seen by backward). Parent F←D.
- E: C (visited), G (new to forward). Parent G←E. G is in backward visited. The searches have met at G.
Rebuild the path. Walk forward parents from G: G←E←C←A, giving A, C, E, G. Walk backward parents from G: G←J, giving J. Joined: A → C → E → G → J, 4 hops.
A plain BFS from A confirms the distance: A is 0; B, C are 1; D, E are 2; F, G are 3; H, J are 4. J is at distance 4. The bidirectional run fetched 6 adjacency lists (A, J, B, C, D, E) and could stop before ever expanding F or the rest of the graph.
Here is the algorithm in Python. It was tested against plain BFS on hundreds of random graphs.
def bidirectional_bfs(graph, source, target):
"""Return a shortest path from source to target, or None."""
if source == target:
return [source]
parent_fwd = {source: None} # visited from the source side
parent_bwd = {target: None} # visited from the target side
frontier_fwd, frontier_bwd = [source], [target]
while frontier_fwd and frontier_bwd:
# Expand the smaller frontier: fewer neighbour lookups.
if len(frontier_fwd) <= len(frontier_bwd):
frontier_fwd, meet = expand(graph, frontier_fwd, parent_fwd, parent_bwd)
else:
frontier_bwd, meet = expand(graph, frontier_bwd, parent_bwd, parent_fwd)
if meet is not None:
return join(meet, parent_fwd, parent_bwd)
return None
def expand(graph, frontier, parents, other_parents):
next_frontier = []
for user in frontier: # one full level at a time
for friend in graph.get(user, ()):
if friend in parents:
continue
parents[friend] = user
if friend in other_parents: # the two searches met
return next_frontier, friend
next_frontier.append(friend)
return next_frontier, None
def join(meet, parent_fwd, parent_bwd):
path, node = [], meet
while node is not None: # meet back to source
path.append(node)
node = parent_fwd[node]
path.reverse()
node = parent_bwd[meet]
while node is not None: # meet forward to target
path.append(node)
node = parent_bwd[node]
return path
Making it work across shards
In production graph.get(user) is a network call, so the code changes in three ways:
- Batch per level. Collect the whole frontier, group IDs by shard, and issue one multi-get per shard in parallel. A level costs one round trip (a few milliseconds), not one per user.
- Cap the work. Stop at
max_hopsand also at a budget, for example 200,000 visited users. Return "not closely connected" rather than melting the cluster. - Skip or sample super-nodes. A celebrity with 50 million followers would explode a frontier. Exclude follower edges from path search, or sample a bounded number of a high-degree node's neighbours. Note that sampling can miss the true shortest path; that is an explicit product trade-off.
Visited sets live in the path service's memory for the duration of one query. A search that touches 100,000 users holds 100,000 × (8-byte ID + parent + hash overhead), a few megabytes per query, which is fine.
Common mistake
Do not run BFS by calling the database once per user. At three levels and thousands of users per level, per-user calls turn a sub-second query into tens of seconds. Always batch a whole level.
Deep dive 2: people you may know
PYMK suggests users you are not yet connected to but probably know. The strongest simple signal is mutual friends: if 12 of your friends are friends with Meera, you probably know Meera.
Friends of friends
For user U:
- Fetch U's friends F₁ … Fₙ.
- For each friend, fetch their friends.
- Count how often each candidate appears. Exclude U and existing friends.
- Rank by count, then by other signals (same school, company, city, contacts uploaded, profile views).
With 200 friends of 200 friends each, that is up to 40,000 candidate occurrences per user. Doing this on every page load for hundreds of millions of users would be a large, repeated cost for a list that changes slowly.
Batch computation
Compute PYMK offline, typically daily, as a distributed job (MapReduce-style or Spark):
Input: edge list (u, v), both directions.
Step 1 For each user f, emit every pair of f's friends (a, b),
meaning "a and b share friend f".
Step 2 Group by (a, b) and count -> mutual_count.
Step 3 Drop pairs that are already friends.
Step 4 For each a, keep the top 50 b by score.
Output: pymk:<a> -> [(b, mutual_count, reasons), ...]
Step 1 is the expensive part: a user with d friends emits d × (d − 1) ordered pairs. A user with 5,000 friends emits about 25 million pairs alone. Mitigations: cap or sample friends for very high-degree users, and skip follower edges of celebrities entirely.
The output is a small list per user, written to a key-value store and served by the suggestion service in milliseconds.
Real-time touches
Pure batch is up to a day stale. Two cheap online adjustments help:
- Remove a suggestion as soon as you connect with, dismiss or block that person (filter at read time against the edge cache).
- Boost right after a new friendship: when you add B, the friends of B become fresh candidates. A stream job can compute B's top friends and merge them into your list.
Deep dive 3: partitioning the graph
Hash sharding by user ID spreads load evenly but ignores structure: your 200 friends are scattered across almost every shard, so each BFS level hits many shards.
Locality-aware partitioning tries to place users who are densely connected on the same shard, for example by region, school or a graph-partitioning algorithm that minimises edges crossing partitions. More edges stay local, so multi-hop queries touch fewer shards.
| Approach | Pros | Cons |
|---|---|---|
| Hash by user ID | Even data and load, simple routing, easy rebalancing with consistent hashing | Neighbours spread everywhere; BFS fans out to many shards |
| Locality / community partitioning | Fewer cross-shard hops for traversal | Uneven shard sizes, hot shards, needs a lookup directory, must re-partition as the graph changes |
| Hybrid: hash for storage, locality in cache | Simple storage, faster traversals for hot regions | Two layouts to maintain |
For this system, hash sharding is the right default because most traffic is one-hop. Mention locality-aware partitioning as an optimisation for the path service if traversal latency becomes the bottleneck.
A directory service (a small, highly cached map from user ID or hash range to shard) lets you move users between shards without changing the hashing scheme. See Scalability for rebalancing strategies.
Deep dive 4: caching hot adjacency lists
Friend lists are read far more than written and are small (about 1.6 KB on average), which makes them ideal cache entries.
- Cache-aside: read cache, on miss read the shard and populate. On write, delete the cached keys for both users.
- Avoid thundering herds: when a popular user's key expires, thousands of requests may miss at once. Use request coalescing (one loader per key) or a short lease so only one request rebuilds the entry.
- Celebrity lists: a list of 50 million followers is not one cache value. Store it in pages (
followers:9:page:0,page:1…) or serve follower counts from a counter and follower pages from the database index. areConnectedchecks: for users with big lists, a Bloom filter (a compact bit array that answers "definitely not present" or "probably present") per user can answer most "no" cases without loading the list.
A cache hit ratio around 90% or better turns about 46,000 average reads per second into a few thousand database reads per second, which is comfortable for dozens of shards.
Scaling and bottlenecks
| Bottleneck | Symptom | Fix |
|---|---|---|
| Hot user (celebrity) | One shard and one cache key overloaded | Replicate hot keys across cache nodes, page large lists, keep follower counts separate |
| BFS fan-out | Path queries slow and expensive | Bidirectional search, batching per level, hop and visit budgets, exclude follower edges |
| Write amplification for mutual edges | Two-shard writes, partial failures | Outbox plus idempotent retries, nightly reconciler |
| Shard growth | Disks fill unevenly | Consistent hashing or directory-based moves; split shards |
| PYMK batch job size | Job runs for hours | Sample high-degree users, incremental computation from change events |
Read replicas absorb list reads; writes go to each shard's primary. Multiple data centres can each hold a full replica of the graph with asynchronous replication, since a friendship appearing a second later in another region is acceptable.
Failure handling
- Second half of a friendship write fails. The outbox row remains, and a worker retries until
B→Ais written. Inserts are idempotent so retries cannot duplicate. - Shard primary dies. A replica is promoted (automatic failover). Writes to that shard fail for seconds; reads continue from replicas, possibly slightly stale.
- Cache cluster loses a node. Requests for those keys miss and fall through to the database. Protect the database with request coalescing and rate limits so a cache loss does not cascade.
- Path query times out on a slow shard. Return a partial "not found within budget" rather than hanging; the feature is not critical.
- PYMK job fails. Keep serving yesterday's lists. Stale suggestions are harmless.
- Lost invalidation event. Give cache entries a TTL (time to live, for example a few hours) so a missed delete eventually self-heals.
See Reliability and recovery for failover patterns.
Trade-offs and alternatives
| Decision | Chosen | Alternative | Why |
|---|---|---|---|
| Storage | Sharded adjacency lists (SQL or wide-column) | Graph database | One-hop reads dominate; mature sharding |
| Mutual edge | Two directed rows | One row with canonical order (min, max) | Two rows make each user's list a single-shard read |
| Shard key | Hash of user ID | Community partitioning | Even load, simple ops; traversals are rare |
| Shortest path | Online bidirectional BFS with budgets | Precomputed distances or landmarks | Pairs are too many to precompute; budgets bound cost |
| PYMK | Daily batch plus real-time filters | Fully online computation | Lists change slowly; batch is far cheaper |
| Consistency between directions | Eventual, with reconciler | Distributed transaction across shards | Simpler, faster; brief asymmetry is acceptable |
An alternative for shortest paths is landmark-based estimation: pick a few hundred well-connected landmark users, precompute everyone's distance to each landmark, and estimate the distance between A and B through the landmarks. It answers "about 3 hops" instantly but does not guarantee the exact shortest path. It is useful when you need only a distance label, not the path.
What interviewers probe
"What if a user has 10 million followers?" Store follower edges in both directions but treat them differently: page the follower list, keep counts in a separate counter, exclude follower edges from path search and PYMK pair generation, and replicate the hot cache key.
"How do you keep the two directions consistent?" Use an outbox with idempotent retries, accept a short window of asymmetry, and run a reconciler that scans for one-sided edges. A two-phase commit across shards would work but adds latency and blocking on failure.
"Why not just use a graph database?" You can for a smaller or more traversal-heavy graph. At 100 billion edges with mostly one-hop reads, sharded adjacency lists plus a cache are simpler to scale and operate.
"How would you make shortest path faster?" Bidirectional search, expand the smaller side, batch per level, cache hot adjacency lists near the path service, and consider locality-aware partitioning to reduce cross-shard calls.
"How do you handle unfriending in PYMK?" Filter suggestions at read time against current edges and a dismissed list; never trust the batch output blindly.
"How do you paginate a friend list that changes during scrolling?" Keyset (cursor) pagination on (created_at, friend_id), so new friends do not shift existing pages.
Interview questions
Q1. Why store a mutual friendship as two directed rows?
So that listing either user's friends is a single lookup on that user's shard. With one canonical row (smaller ID first), listing friends would need two indexes or a scatter query across shards. The price is two writes per friendship and a short consistency window.
Q2. Explain bidirectional BFS and why it is faster.
You run BFS from both the source and the target, always expanding the smaller frontier, and stop when the frontiers meet. A path of length L is found after each side explores about L/2 hops, so work drops from roughly bᴸ to 2·b^(L/2). With degree 200 and a 4-hop path, that is 80,000 lookups instead of 1.6 billion.
Q3. Does stopping at the first meeting node guarantee a shortest path?
Yes, when each side expands a full level at a time and you check for a meeting as each node is discovered. Any shorter path would have caused a meeting at an earlier level. If you mixed partial levels or weighted edges, you would need more care (for weighted graphs, bidirectional Dijkstra has a different stopping rule).
Q4. How do you run BFS when adjacency lists live on hundreds of shards?
Treat each level as one batch: group the frontier's IDs by shard and send parallel multi-gets. Add budgets on hops and visited nodes, and use a cache in front of the shards. That keeps a query to a handful of round trips.
Q5. How would you compute "people you may know"?
Count mutual friends by generating, for each user, all pairs of their friends; group by pair and count; drop existing friends; keep the top N per user. Run it as a daily batch job, store results in a key-value store, and filter at read time against new friendships and dismissals.
Q6. What is the main risk in the PYMK batch job?
Quadratic pair generation for high-degree users: a user with d friends emits about d² pairs. Celebrities can dominate the job's runtime and memory. Sample or cap their friends and exclude follower edges.
Q7. Graph database or sharded key-value store?
Choose by access pattern. If most queries are one hop and the graph is enormous, sharded adjacency lists are simpler and faster. If queries are deep, varied traversals on a graph that fits a cluster comfortably, a graph database saves a lot of application code.
Q8. How do you shard the graph?
Hash the user ID (or use consistent hashing) and keep each user's adjacency list on one shard. A directory service can map users to shards for moves. Locality-aware partitioning is an optimisation for traversal-heavy workloads.
Q9. How do you cache friend lists, and how do you invalidate them?
Cache-aside keyed by user ID. On any edge change, delete both users' keys after the database commit, driven by change events, and set a TTL as a safety net. Large lists are paged rather than stored as one value.
Q10. A friendship write succeeded on A's shard but failed on B's. What now?
The outbox record on A's shard drives retries of B→A until it succeeds; the insert is idempotent. Meanwhile A sees B but B does not see A. A periodic reconciler finds and repairs any remaining one-sided edges.
Q11. How do you answer areConnected(a, b) quickly?
Point lookup on primary key (a, friend, b) on A's shard, usually served from the cached adjacency set. For users with very large lists, a per-user Bloom filter quickly rules out most negatives.
Q12. How would you estimate degrees of separation for every pair?
You would not; there are far too many pairs (500M² is about 2.5 × 10¹⁷). Compute on demand with bounded bidirectional BFS, or precompute distances to a few hundred landmark users and estimate through them when an approximate number is enough.
Q13. What consistency model do you choose and why?
Eventual consistency across the two directions and across regions, with read-your-writes for the user who made the change (route their next read to the primary or update their cache entry). Social relationships tolerate a second of lag; availability matters more.
Key takeaways
- A social graph is nodes (users) plus edges (relationships); most traffic is one-hop reads such as "list friends" and "are A and B connected".
- Store adjacency lists sharded by user ID, and write mutual friendships as two directed rows so each list is a single-shard read.
- BFS cost grows as degree to the power of hops; bidirectional BFS halves the exponent and is the core of degrees-of-separation.
- In a distributed graph, batch each BFS level per shard and enforce hop and visited-node budgets.
- People you may know is a mutual-friend count computed in batch, served from a key-value store, and filtered in real time.
- Celebrities break naive designs: page their lists, separate counters, and exclude follower edges from expensive computations.
- Cache adjacency lists with delete-on-write invalidation and TTLs; the cache absorbs most read traffic.
- Choose a graph database only when deep traversal dominates; at huge scale, sharded storage plus application-level traversal is the common choice.
Next lesson
Continue with Design a sales ranking system.

