Scaling Facebook’s Social Graph with TAO: Cache‑First Reads and Eventual Consistency
Learn how Facebook’s TAO store uses a cache‑first approach to serve billions of social‑graph edges with low latency, and what trade‑offs arise from its eventual consistency model.
10 May 2026, 03:12 UTC

Why Facebook Needed a New Data Store
At peak usage, Facebook serves billions of friend‑graph edges per second. A traditional relational database cannot deliver the sub‑millisecond reads required for news‑feed generation while also handling the constant stream of friend‑request writes.
TAO’s Cache‑First Design
TAO layers an in‑memory cache in front of a sharded MySQL store. Reads are served from the cache whenever possible; a miss falls back to MySQL and populates the cache. Writes update the cache synchronously and are persisted to MySQL asynchronously.
Sharding and Lease‑Based Locking
Objects are hashed by their 64‑bit ID to determine a shard. Each shard runs a lease‑based lock manager that grants short‑term leases to replicas, allowing concurrent writes while preventing split‑brain scenarios.
Read and Write Paths in Detail
Read path: client → TAO client library → local cache → (if miss) storage node → cache fill → response.
Write path: client → TAO client library → cache update → async queue → MySQL replica → durability.
Worked Example: Fetching a Friend List
# Pseudo‑code using TAO client (language‑agnostic)
start = now()
friends = tao.get_edges(user_id, 'friend')
latency = now() - start
print(f'Retrieved {len(friends)} friends in {latency:.3f} ms')
When the cache holds the user’s edge list, the call returns in under 1 ms. After a new friend edge is written, the asynchronous propagation means the updated list appears on subsequent reads once the cache is refreshed—typically within the replication window (tens of milliseconds to a few seconds depending on load).
Trade‑Offs and Limitations
TAO favors eventual consistency. A newly created friendship may not be visible to all readers immediately, which can affect features that rely on real‑time updates (e.g., 'just added' badges). Operators must also manage cache warming, failure detection, and cross‑region replication, which adds operational overhead compared to a simple DB‑only setup.
Practical way to check the impact: enable TAO’s statistics endpoint (e.g., http://tao‑host:port/stats) and watch the cache_hit_ratio metric. A dropping ratio signals that the cache is not keeping pace with write traffic, and latency measurements will rise accordingly.
Actionable Takeaways
- Instrument your TAO client to log read latency and cache‑hit/miss counters.
- Set alerts on
cache_hit_ratiofalling below a threshold (e.g., 95 %) to trigger cache‑warming jobs. - When building real‑time features, consider reading from the storage layer directly or using a version‑vector to detect stale cache entries.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.