Opening Facebook feels like loading one page. Behind that page are several different questions: which people are your friends, what have they posted, which comments belong to each post, and what are you allowed to see?
Those questions explain more about the storage design than a list of database products. Facebook needs to retrieve connected pieces of information quickly while people continually create and change those connections. A successful design must make the familiar actions cheap without turning every page view into a search across the entire network.
Introduction
Facebook's engineering team has published detailed accounts of TAO, its social graph storage service. The original description dates from 2013, so it is a documented architectural case study rather than a complete inventory of Facebook's infrastructure today.
The useful idea is that a graph-shaped product can run on relational storage. We will look at that documented foundation, then use a simplified social application to explain the decisions it creates. The example records and request flows below are teaching models, not Facebook's internal schemas.
For someone learning system design, this distinction matters. Understanding why related data is grouped, cached and replicated is more useful than memorising that a particular company once used a particular database.
Objects and Relationships Form the Social Graph
A graph contains objects, often called nodes, and connections between them, called edges. A person and a post are objects. Authorship, friendship and liking a post are relationships.
Meta's account of its Dragon query engine describes this model as objects and associations. It explains that TAO was designed for high volumes of single-hop queries: starting with one object and retrieving its immediate connections. Meta's social graph explanation
Imagine Alice writes post P and Ben likes it. The post record holds the content, while relationships connect Alice to P and Ben to P. Asking who liked P does not require searching the text of every post. It requires finding the relevant connections.
Direction is important. “Posts written by Alice” and “author of post P” are opposite lookups. A small application may support both through indexes on one table. A distributed service may maintain separate representations when that makes important reads cheaper. Either way, the query must influence the model.
Relational Storage Underneath a Graph Interface
The published TAO design uses MySQL for durable storage, with a service interface focused on objects and associations. Caching clusters handle reads and coordinate writes and cache consistency. TAO's published architecture
This separates the product's data model from its storage engine. Developers can work with graph concepts without sending arbitrary graph traversals directly to a specialist graph database.
Our smaller application could begin with three ordinary relational tables:
- Users: an ID, display name and account state.
- Posts: an ID, author ID, body, visibility and creation time.
- Friendships: two user IDs and the state of their relationship.
Indexes would support the operations the application actually needs. An index beginning with author ID can help retrieve that author's posts. A friendship lookup needs to find the pair of users efficiently.
The service might expose getRecentPosts and getFriends operations. That boundary gives the product a stable interface even when the physical tables or partitioning change. It also centralises validation instead of asking every feature team to remember the storage rules.
Sharding Spreads the Storage Work
TAO's published design partitions objects and associations into many shards backed by MySQL databases. A shard is a portion of the overall dataset, allowing storage and requests to be distributed across machines.
In our application, putting related records together can make common requests cheaper. Fetching a person's recent posts should ideally reach a small, predictable number of storage locations. A request that contacts every shard becomes more expensive as the system grows.
There is a tradeoff. A famous account can receive far more traffic than an ordinary account. Evenly distributing record counts does not evenly distribute requests. One shard might hold little data but serve a disproportionate share of the traffic.
It helps to distinguish logical shards from physical servers. If records belong to a stable logical partition, that partition can be moved to another machine without changing every record's identity. Routing information then answers which server currently owns the partition. This adds operational complexity, but makes capacity changes more manageable.
Caching Makes Repeated Reads Affordable
Consider a popular post being opened by thousands of people. Most requests need the same author details and post content. Retrieving those values from memory can avoid repeatedly reading the durable database.
The hard part is deciding what happens after an edit. Changing the stored post while leaving its old cached copy indefinitely would produce inconsistent pages. A cache strategy needs an update or invalidation mechanism, plus a way to recover when notifications are missed.
For our application, a cache key might identify post P and hold its latest known content. The cache is an optimisation: the database remains the authoritative record. Losing the cached entry should cause a slower read, not permanent loss of the post.
A cache outage can still threaten availability. If every request suddenly falls through to the database, normal traffic becomes an unexpected surge in database work. Request coalescing, concurrency limits and gradual cache warming are possible protections. They solve different problems, so the team should measure which failure is actually occurring before adding them.
Following a New Post Through the System
Suppose Alice submits “First day at the new job.” In a simplified implementation, the request first authenticates Alice and validates the content and audience selection. The service then saves the post and its author relationship.
The acknowledgement should have a clear meaning. Does it mean the post has reached durable storage, or only that a server accepted the request into memory? Users naturally interpret “posted” as something that will survive closing the browser.
After the core write succeeds, background work can prepare notifications, update search or generate feed candidates. Those tasks do not all need to delay the initial response.
However, “save the post, then send an event” creates a failure gap: the process could crash between those actions. In a small relational design, a transactional outbox can store the post and a pending event in the same transaction. A worker later delivers that event. This is a possible design for our example, not a claim that TAO uses an outbox.
A Post and Its Audience Are Different Concerns
Storing a post successfully does not determine who should see it. A social application must combine content with audience rules, account relationships and the viewer's permissions.
A cached post body might be reusable across many viewers, while “can this person read it?” depends on the current viewer. Treating those answers as one universally reusable cache entry risks returning private information to the wrong person.
Consider Alice changing a post from public to friends-only. A previously generated feed candidate may still reference the post. The read path needs to check the current access rules rather than assuming the old candidate list grants access forever.
The same issue affects search results, notifications and previews. Removing an item from one page does not necessarily remove every copy or reference. Designing the lifecycle of derived data is part of implementing privacy correctly, even when each individual storage operation looks simple.
Feed Ranking Is Separate from Durable Content
The social graph tells a system which objects are related. It does not, by itself, decide which post deserves the first position in someone's feed.
A useful conceptual pipeline is to collect candidate post IDs, apply visibility rules, rank candidates, and retrieve the content needed for display. These stages have different storage and computation needs.
Keeping candidate IDs is much smaller than duplicating every post body for every reader. It also lets the display stage retrieve an updated post. The cost is additional lookups and the need to cope with missing or inaccessible records.
There are also choices about when to prepare candidates. Preparing work when an author posts can make later reads cheap, but becomes expensive for an author with many followers. Preparing everything when a reader opens the feed avoids that immediate fan-out, but increases read-time work. A teaching design can combine approaches, provided it defines which authors or audiences use each path.
Replication Requires Clear Consistency Choices
The original TAO description uses a primary region per shard, replication to other regions and eventual consistency as the default. Different copies can temporarily disagree while updates propagate.
In our example, Ben might see a new like before another reader sees the updated count. That may be acceptable. Showing content after access has been revoked is a different problem and deserves a stricter policy.
Read-your-writes behaviour is another useful requirement: after Alice edits her post, she should normally see the updated text. Otherwise, she may retry or assume the edit failed. A system can provide this experience through request routing or version-aware reads without requiring every unrelated operation worldwide to use the same consistency level.
Replication is also different from backup. A mistaken deletion may quickly reach every replica. Recovering from that mistake requires retained history or backups and a tested restoration process. Extra copies improve resilience only when their failure modes are understood.
Retries Must Not Create Duplicate Actions
Imagine Alice submits a post, the server commits it, and the response is lost. Her browser cannot tell whether the operation failed or merely the acknowledgement failed.
If retrying creates a new post each time, one click can become several identical posts. A possible solution is a client-generated request identifier that the service remembers alongside the result. Repeating that identifier returns the earlier outcome rather than creating another post.
The storage operation must enforce uniqueness atomically. Checking for a request ID and then inserting in a separate unprotected step leaves a race when two retries arrive together.
This is a general distributed-systems lesson: a timeout describes what the caller knows, not necessarily what happened in storage. Any user action that crosses a network needs an explicit retry policy.
What a Smaller Application Should Build First
A new social product rarely needs a global graph service immediately. A relational database with suitable indexes, bounded queries and clear permission checks can support a useful first version.
Start by measuring the expensive operations. Is the feed slow because it scans too many rows, performs one query per post, or waits on an external service? These problems require different fixes. Adding a cache before understanding the query can hide an inefficient design until the cache misses.
Useful checks include loading a feed with a cold cache, editing a popular post during heavy reads, retrying a timed-out submission and revoking access while old references still exist. These reveal whether the product remains correct under realistic conditions, not just whether a benchmark returns quickly.
Add partitioning when a measured capacity or isolation problem justifies its operational cost. The ability to run many databases also means learning to migrate, monitor, repair and restore many databases.
Watch the Shape of a Request
A useful performance exercise is to count the storage calls needed to render one feed page. Ten posts should not automatically mean ten separate author queries, ten permission queries and ten relationship queries performed sequentially.
Batch compatible lookups and set a budget for how much fan-out a request can create. Then measure that budget with realistic data, including missing records and cold caches. This makes a slow page easier to explain: the team can distinguish too many dependent calls from an individually slow database operation. Optimisation becomes a response to evidence rather than a guess about which component sounds expensive.
The Big Picture
Facebook's published architecture demonstrates how a graph interface, relational persistence, partitioning and caching can work together. Each addresses a different part of a demanding social workload.
The most transferable lesson is to begin with relationships and access patterns. Decide which records need to stay together, which answers can be reused, and which changes must become visible promptly. Then walk through retries, permission changes and partial failures.
A storage design is successful when those ordinary product actions remain understandable and reliable as the workload grows. The database brand is only one part of that design.
