Storing and querying graphs
Social graphs and other graph data in system design: adjacency lists in a key-value store, graph databases, friends-of-friends and mutual friends, celebrity nodes, sharding a graph, and what Facebook's TAO does.
Reading is half of it. See this used in a real interview: walk through Design a News Feed →
Followers, friends, likes, group memberships and "people you may know" are all graphs: things (nodes) connected by relationships (edges). Most system design questions with a social component need a graph store, and interviewers probe how you store it, how fast common queries are, and what happens to the account with 300 million followers.
Most queries are one hop
The queries real products run at scale are usually simple:
- Who does A follow? Who follows A?
- Does A follow B?
- How many followers does A have?
- Which of A’s friends liked this post?
These are one-hop lookups from a known node. Deep traversals ("shortest path between any two people") are rare in products and are usually computed offline. Design for one-hop first.
Adjacency lists in a key-value or wide-column store
The workhorse design: store each node’s edges as a list keyed by the node, in both directions.
following: (user_id, followee_id) → created_at partition by user_id
followers: (user_id, follower_id) → created_at partition by user_id
counts: user_id → { followers, following }- "Who does A follow" is one partition read; "who follows A" is one partition read in the other table.
- "Does A follow B" is a single key lookup.
- Writes update both directions (and the counts), ideally in one transaction or through an idempotent asynchronous job that repairs the second direction.
This scales horizontally in Cassandra, DynamoDB or sharded MySQL and is what most large social products use under the hood. See data modelling for reads.
Graph databases
Neo4j, Amazon Neptune and similar systems store nodes and edges natively and run traversal queries (Cypher, Gremlin) efficiently, following pointers rather than doing joins.
Good for: multi-hop queries over moderately sized graphs (fraud rings, knowledge graphs, permission hierarchies, recommendations computed on demand), and rich queries over relationships.
Less good for: web-scale social graphs with billions of edges and very high read rates, because traversals cross shard boundaries and partitioning a graph well is hard. Large social networks generally use sharded adjacency lists with a caching layer instead.
Facebook’s TAO
Facebook’s published design is worth knowing as a reference: objects (users, posts) and associations (typed edges such as likes, friend) stored in sharded MySQL, with a large, write-through, geographically distributed cache in front that answers almost all reads. Association lists are kept in time order with counts, so "the 50 most recent likes of this post" and "how many likes" are fast. It is essentially adjacency lists plus aggressive caching, tuned for one-hop reads at enormous scale.
Friends of friends and mutual friends
- Mutual friends of A and B: intersect two adjacency lists. With lists of a few hundred, that is trivial in memory.
- Friends of friends ("people you may know"): A has 300 friends, each with 300 friends, so 90,000 candidates per user, too many to compute per request for every user. It is computed offline (batch graph jobs) or incrementally, ranked, and stored as a recommendation list per user.
Celebrity nodes
Degree is extremely skewed: most users have hundreds of followers; some have hundreds of millions.
- Storage: a single partition holding 300 M followers is too large; split huge adjacency lists into sub-partitions (user_id, bucket).
- Reads: "list all followers" is never done synchronously for a celebrity; it is a paginated or batch operation.
- Fan-out: pushing a celebrity’s post to every follower is the classic problem solved by fan-out on read for high-degree nodes. See Design a News Feed.
- Counts: follower counts on celebrities change constantly; use sharded counters.
Sharding a graph
Edges connect nodes on different shards, so some queries always cross shards. The usual approach is to shard by the source node (all of A’s edges live with A), which makes one-hop queries single-shard and accepts that two-hop queries fan out. Smarter partitioning that keeps communities together (graph partitioning) reduces cross-shard traffic but must be recomputed as the graph changes; few systems do it online.
Caching
Adjacency lists for active users are read far more than they change, so they are cached heavily, with the edge write invalidating or updating both endpoints’ cached lists. "Does A follow B" checks for the people a user sees most often can be served from a per-user cached set or a Bloom filter. See caching.
In the interview
Store edges as adjacency lists in both directions, keyed by the source node, with counts maintained alongside. Say which queries are one hop and cheap, which are computed offline, and how high-degree nodes are handled. Mention a graph database only if the product needs multi-hop queries on demand.
Checklist
- Edge tables in both directions, partitioned by source node.
- Counts maintained separately.
- One-hop queries online; multi-hop queries offline.
- High-degree nodes: split lists, no synchronous full scans, fan-out on read.
- Caching of adjacency lists and membership checks.
- Graph database only when traversals are the product.