M3 Design Social Graph (TAO)
Loading learning experience...
Lecture transcript
Read the narration for M3: Design Social Graph (TAO)
Design Social Graph (TAO)
Welcome everyone, today we will explore how TAO builds fast social graph reads using objects and associations, with cache first design, smart sharding, and read after write consistency.
From distributed cache to a social graph
Dr. Wei: Today we are bridging two ideas: the distributed caching patterns we studied last time, and the problem of serving a huge social graph with low latency and high availability.
Sam: So we are taking the caching tools we learned and applying them to social graph reads, where tiny delays add up and availability matters?
Dr. Wei: A social graph workload is read heavy, highly skewed, and sensitive to tail latency. That combination makes caching not just an optimization, but a core part of the system design.
Why a custom graph store?
Dr. Wei: When you build a social graph at massive scale, the storage choice is not just a database preference; it directly determines whether common actions feel instant or sluggish.
Sam: If we tried to do this with just relational tables, is the core problem that joins become too slow and too expensive as the relationship table grows?
Dr. Wei: One option is to model relationships in relational tables, but then answering everyday questions like friends of friends often turns into large join-heavy queries that grow expensive as the data grows.
Dr. Wei: Another option is an off-the-shelf graph database, but a general system may not align with our exact access patterns, caching strategy, and strict latency targets for reading and updating the social graph.
Data model: objects and associations
Dr. Wei: Before we talk about storage or APIs, we need a simple data model for a social graph. The core idea is that most social products can be described using just two primitives, and everything else is built by combining them.
Sam: Two primitives meaning a thing, and a relationship between things, right? Like users and the links between them.
Dr. Wei: The first primitive is an object: an entity with an identifier and attributes. Think of objects like users, photos, pages, and posts, where the attributes capture details like a name, timestamp, caption, or privacy setting.
TAO API: small set, huge leverage
Dr. Wei: In TAO, we aim for a deliberately small public interface that still supports a very large range of social product behavior. The core idea is leverage: if the primitives are stable and expressive, many different features can be built on top without adding new special cases. Concretely, TAO centers on a small set of operations like getObject for object reads, assocGet and assocRange for relationship reads, assocCount for fast counts, and assocAdd and assocDel for relationship writes.
Sam: In an interview, do you want me to name specific operations, like getObject, assocRange, and assocGet, rather than describing a huge set of feature endpoints?
Dr. Wei: Yes. Instead of exposing dozens of endpoints for every feature, TAO focuses on a few consistent operations that cover reading objects with getObject, reading edges with assocGet and assocRange, sizing results with assocCount, and updating relationships with assocAdd and assocDel. That keeps the system easier to reason about, safer to evolve, and easier for many teams to use correctly.
Two-level architecture: TAO cache $+$ MySQL
Dr. Wei: Now that we have a graph service that needs to answer lots of small, repetitive lookups quickly, we can talk about the simplest practical backend shape: a fast serving layer in front of durable storage.
Sam: So the intent is a cache-first read path, with a database behind it as the source of truth, and the hard part is keeping them from drifting in ways users notice?
Dr. Wei: Reads try the cache first because most requests repeat across users and time, so serving from memory cuts latency and protects the databases. When there is a miss, TAO reads from MySQL, then fills the cache so the next request is fast. Writes go to MySQL for durability, and TAO invalidates or updates related cache entries so future reads do not serve stale data.
Graph caching: cache nodes and edge lists
Dr. Wei: Let’s make caching concrete in a social graph system. We have to decide what we cache, what we key it by, and what ends up being requested so often that it becomes the hot path. Before we get specific, I want you to reason from first principles: what pieces of a graph get read repeatedly in a product, and what does that imply for caching?
Sam: I’d guess we cache the user or object nodes themselves, like the profile data, because lots of pages need that. And maybe we cache the lists of connections, like friend lists or followers, since those drive feeds and mutual connections.
Dr. Wei: Good. Now make a falsifiable prediction: if we can cache both nodes and edge lists, which cache keys do you expect to be the hottest, and why? For example, do you think the hottest keys are a single object’s node record, a single object’s outgoing edge list for a given edge type, or something else? Commit to one and give your reasoning.
Sam: I’ll commit to edge lists being the hottest, especially follower or friend lists for popular accounts, because so many features need to read them repeatedly and they change less often than the feed itself.
Dr. Wei: Now let’s do the compare moment. In practice, the hottest cache keys are usually specific association lists for a small set of very high degree objects, like the followers list for a celebrity or the memberships list for a huge group, often for one edge type at a time. Node records get hits too, but degree skew means a tiny fraction of objects dominate reads, and their edge lists are fetched over and over for fan out and mutuals.
$Shard-by-id1$: optimize outgoing edges
Dr. Wei: Before we pick a sharding key, we need to be clear about what the system does most of the time. In a social graph service, a very common request is: given a source node, fetch its outgoing edges, essentially an adjacency list. At the same time, reverse lookups, like finding all incoming edges for a target, can be just as important in some products and for some objects.
Sam: If we shard by the source, we optimize the common adjacency list reads, but then anything like who follows this account becomes harder unless we add an index or a separate lookup path, right?
Dr. Wei: That leads directly to the choice on this slide: shard by id one, the source identifier, to optimize outgoing edges. The tradeoff is that reverse lookups by the destination id may need extra work, and for high profile targets those incoming edge queries can be very hot. Common mitigations are to materialize reverse edges with a dual write to a reverse table, maintain a secondary index, or route follower style reads to a separate lookup service built for that access pattern.
Read-after-write: what users psychologically demand
Dr. Wei: When people use a social app, they carry a simple expectation: if I change something, I should immediately see that change reflected when I look again. That expectation is read-after-write consistency, and it is more psychological than technical.
Sam: So if I post, like, or follow someone, I expect my next refresh to show it right away, even if the system is distributed?
Dr. Wei: Exactly. In TAO-style social graphs, the challenge is that data may be cached and replicated, so different machines can temporarily disagree. The contract we want is: after a user performs a write, reads from that same user should not appear to go backward in time. Caches matter here because they must avoid serving an older version right after the user just wrote a newer one.
Association queries: building blocks for features
Dr. Wei: On this slide, we will make the idea of association queries feel concrete. In a social graph system like TAO, a lot of product features reduce to one basic need: given a person, fetch a set of related entities quickly and predictably.
Sam: So the reusable building block is basically: from this person, get me their connections of a certain type, maybe with a limit and ordering, and then features layer on filters and ranking?
Dr. Wei: We will use the friend list as the running example: you ask for a person’s friends, usually with a limit, an ordering, and sometimes filters like only active users. The big takeaway is that association queries are the reusable building blocks that higher level features compose into timelines, recommendations, and privacy aware views.
Graph traversals: friend-of-friend is not one query
Dr. Wei: Today we are looking at why a friend-of-friend request is not a single query in a system like TAO. The key idea is that graph traversals tend to multiply work across hops, and that multiplication shows up as both extra calls and extra latency at the tail.
Dr. Wei: Quick estimate question: suppose a typical user has about 500 friends. If we want two-hop candidates, meaning friends of those friends, how many primitive calls do you think we make, and how much data do we end up touching?
Sam: We probably do one call to fetch the first hop, then around 500 more calls to fetch each friend’s friend list. Data-wise, we might read 500 friend lists, so on the order of 500 times 500 edges, like roughly 250,000 edge reads, before we even deduplicate and rank.
Dr. Wei: Exactly. The call count and the amount of touched data blow up multiplicatively with fanout, and that amplifies the 99th percentile latency because you are now waiting on many requests, and the slowest ones dominate. That is why TAO exposes small primitives and you build traversals by repeating them, but in practice you also add limits, early stopping, caching, and sometimes offline precompute so the online path does not try to expand the full two-hop neighborhood.
Failure modes interviewers expect you to name
Dr. Wei: In interviews, naming failure modes is less about listing everything that can go wrong and more about showing a repeatable way to reason under pressure. We will organize the common failures you should mention into a few buckets, and we will use one template each time: symptom, likely root cause, and a practical mitigation.
Sam: When I list failure modes, should I prioritize ones that affect user-visible correctness, like stale edges, before pure performance issues, or do you want a balanced list?
Dr. Wei: First bucket: correctness and consistency. The symptom is stale edges, missing edges, duplicate edges, or a follower count that does not match what the edge list implies. The root cause is write loss, non atomic multi write updates, cache incoherence, or eventual consistency windows that are not accounted for. Mitigations include idempotent writes, write ahead logging or durable queues, cache invalidation tied to the write path, background reconciliation, and being explicit about read your writes expectations.
Dr. Wei: Second bucket: load and hot-spotting. The symptom is timeouts or 99th percentile latency spikes concentrated on a few shards or a few high degree users. The root cause is skewed key distribution or uneven fan out. Mitigations include better sharding keys, splitting hot partitions, adding caching for hot reads, and rate limiting or backpressure for heavy writers.
Tradeoffs: what level is your answer?
Dr. Wei: Before we wrap up, let’s talk about what level of answer you should give on a Design Social Graph, TAO-style question. A strong answer is not only correct, it is calibrated to the time you have, the signal the interviewer wants, and the risks in the system.
Sam: If I only have, say, fifteen minutes, should I aim for the mid-depth answer with clear read and write flows plus caching and consistency, and only add failure modes if the interviewer asks?
Dr. Wei: Your goal is to match depth to what is being tested. If the prompt is about news feed performance, spend your time on read path tradeoffs and hot key mitigation. If it is about relationship updates, focus on write amplification, propagation, and stale data. If the interviewer asks for reliability, talk about what breaks, how you detect it, and how you recover.
Exit ticket: reason about a write and the next read
Dr. Wei: Let’s end with a quick exit ticket about a very common social graph situation: someone updates their friend state, and then we immediately do a read. Your job is to reason about what the reader might observe and why.
Sam: Am I expected to call out both cache invalidation timing and database replication lag as reasons the read might briefly show the old state?
Dr. Wei: In your explanation, be specific about where the read is served from. Consider whether you are reading from a cache versus a backing store, and whether anything needs to propagate or be invalidated before the freshest result is visible. One or two clear sentences is enough, as long as you justify what could happen.
Thank you for watching!
Thanks for watching. Subscribe and share if you found this useful—see you next time!