Scaling the Social Graph: How Facebook's TAO Manages Distributed Reads
Explore how Facebook's TAO handles billions of social graph edges using a read-through cache and write-behind MySQL sharding to achieve sub-millisecond latency.
03 Jun 2026, 19:18 UTC

The problem: The latency cost of a trillion-edge graph
When a user loads their news feed, the system must resolve dozens of relationships: who are their friends, what pages do they follow, and who liked those specific posts? Performing these joins across a traditional relational database at the scale of billions of users would create catastrophic latency and database contention.
The core challenge is balancing the durability of a disk-based store with the sub-millisecond response times required for a fluid user interface. The solution is TAO (The Associations and Objects), a distributed graph service that decouples the logical graph from the physical storage.
The TAO Architecture: Objects and Associations
TAO treats the social graph as a collection of Objects (nodes like Users, Pages, or Posts) and Associations (edges like "friend of" or "authored by"). These are stored as rows in MySQL shards, but the application rarely interacts with MySQL directly.
To maintain performance, TAO implements a read-through caching layer. The request flow follows a strict priority:
- Cache Hit: The request hits the TAO cache (a memcached-style layer). If the object or association is present, it is returned immediately.
- Cache Miss: TAO queries the specific MySQL shard owning that data, populates the cache with the result, and then returns it to the client.
By partitioning data by object ID across thousands of shards, TAO ensures that no single database server becomes a bottleneck for the entire graph.
The Write Path and the Consistency Gap
To keep write latency low, TAO employs a write-behind model. Unlike the read path, writes do not stop at the cache.
- The write request is sent directly to the authoritative MySQL shard.
- MySQL updates the record synchronously to ensure durability.
- The update is then propagated asynchronously to the cache layer.
This design prioritizes write availability and speed over immediate consistency. Because the cache update happens after the database commit, there is a brief window where a read request might return stale data from the cache before the asynchronous update arrives.
Worked Example: Retrieving an Association
In a production environment, services interact with TAO via Thrift (a cross-language RPC framework). Below is a conceptual implementation of fetching a list of associations (e.g., a user's friends).
// Requires: fbthrift client library
// Permissions: tao.read access to the 'social_graph' namespace
import facebook.tao.TaoService;
import facebook.tao.types.TaoGetResult;
TaoService.Client client = new TaoService.Client(protocol);
// The ID of the user whose friends we are fetching
const int64 USER_ID = 1234567890;
try {
// Requesting the 'friends' association for the given USER_ID
TaoGetResult result = client.get(USER_ID, "friends", null);
if (result.isCached()) {
// Path: Client -> TAO Cache (Latency: ~microseconds)
log.info("Served from cache");
} else {
// Path: Client -> TAO -> MySQL -> Cache (Latency: ~milliseconds)
log.info("Served from MySQL shard");
}
List friendIds = result.value().getListOfInt64();
processFriends(friendIds);
} catch (TaoException e) {
// Handle shard unavailability or network timeouts
log.error("TAO request failed: " + e.getMessage());
fallbackToDefaultView();
}
Execution Context: This code runs within an internal service using a Thrift-generated client. It requires a service account with specific namespace permissions.
Verification: Success is verified by monitoring the isCached() boolean. A healthy system should maintain a cache hit ratio typically above 95% for hot objects.
Risks: A "cache stampede" can occur if a highly popular object expires from the cache, causing thousands of simultaneous requests to hit a single MySQL shard. TAO mitigates this via request coalescing and load-shedding.
Trade-offs: Eventual Consistency
The primary limitation of TAO is eventual consistency. If a user updates their profile and immediately refreshes the page, they might see the old data if the asynchronous cache update hasn't completed.
For critical features (like security settings or notification badges), Facebook uses specific mitigations:
- Read-after-write tokens: The write operation returns a version token. The subsequent read includes this token, and TAO ensures the cache is updated to at least that version before responding.
- Strongly consistent reads: Bypassing the cache to read directly from the MySQL master for a specific, high-importance request.
Actionable Summary
When building a system to handle massive graph-like data with low latency, apply these TAO-inspired principles:
- Shard by ID: Distribute your data across storage nodes based on a unique identifier to prevent hotspots.
- Read-Through, Write-Behind: Use a cache for all reads, but write synchronously to the database and asynchronously to the cache to maintain high write throughput.
- Monitor Hit Ratios: Track the percentage of requests served by the cache; a drop here is a leading indicator of database stress.
- Plan for Staleness: Identify which parts of your UI can tolerate eventual consistency and which require a "strong read" path.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.