M2 Design Messenger
Loading learning experience...
Lecture transcript
Read the narration for M2: Design Messenger
From News Feed to Messenger: different latency, different guarantees
Dr. Wei: Today we pivot from News Feed to Messenger, and the key idea is that the product goal changes the system design.
Sam: So even if both are just showing content, the difference is that feed can be a bit late, but chat feels broken if it is late?
Dr. Wei: Exactly. A feed can tolerate seconds of end-to-end delay if it improves relevance, reduces load, and gives a stable experience when you scroll. Messaging is different: it feels broken if the round trip is slow, if delivery is uncertain, or if people disagree about what happened first, so we need tighter latency and clearer delivery and ordering guarantees at massive scale.
Why messaging is deceptively hard
Dr. Wei: Messaging sounds simple because it feels like a chat bubble should just appear instantly. But the real promise users expect is both speed and reliability at the same time, and those goals often pull the system in opposite directions. That tension is why messaging systems are deceptively hard to design.
Sam: If we just use at most once delivery for speed, users probably will not notice, right? Missed messages seem rare.
Dr. Wei: They notice the first time it happens, and trust is hard to win back. Now look at the “correctness traps” line: reordering and duplicates happen naturally with retries and multi-device; split-brain happens if two servers both think they are in charge. The implication is we need idempotency to remove duplicates, and a single authority per ordering domain to prevent conflicting histories.
Scope and success metrics
Dr. Wei: Start with the in-scope bullet: one to one, groups, presence, receipts, and typing. The implication is we are committing to both message delivery and a stream of small real time signals, which means the system must handle many tiny events, not just big message writes.
Sam: For the SLOs, why not just say one second for p ninety nine? That seems easier and still feels fast.
Dr. Wei: The SLO bullet says send to deliver under one hundred milliseconds at p ninety nine, plus durable storage. The tradeoff is cost: tighter latency needs persistent connections and regional proximity, and durability needs an append before we claim success. Finally, look at de-scope: search over end to end encryption and spam and abuse machine learning. The implication is we avoid features that require heavy server side indexing or complex classifiers in v1, so the core delivery path stays predictable.
Back-of-the-envelope scale
Dr. Wei: Before we reveal any numbers, let’s do quick back-of-the-envelope estimates for three things we’ll need to size: the write path, the storage we ingest per day, and the presence heartbeats.
Dr. Wei: First, traffic: if we have about one hundred billion messages per day, what do you expect that to be in messages per second on average, and what might peak look like? If you divide by about eighty six thousand seconds per day, you land around 1.2 million messages per second on average, and a reasonable peak target is about 3 to 5 million messages per second. Design implication: you size partitions and throughput targets for multi-million writes per second, not a single writer.
Sam: For payload, can we ignore metadata and just count the text? Most messages are short.
Dr. Wei: Let’s predict it: even if the text averages only about one hundred bytes, what do you think metadata adds per message for ids, timestamps, receipts, and encryption overhead? In practice it is often around a thousand bytes of overhead, so you end up near 1.1 kilobytes per message. Multiply that by one hundred billion messages per day and you get about one hundred ten terabytes per day before replication. Design implication: storage cost is driven by retention and replication, so plan tiering and compaction early.
Dr. Wei: Finally, presence: quick estimate—if you have about five hundred million concurrent connections and each sends a heartbeat every thirty seconds, what heartbeat rate do you expect? It is five hundred million divided by thirty, which is about seventeen million heartbeats per second. Design implication: keep per-connection state and per-heartbeat work extremely cheap, or the heartbeat fleet dominates your infrastructure.
APIs and client events
Dr. Wei: In a messenger, the client and server need a clear contract for three things: writing messages, reading and catching up, and staying updated in real time. A key design choice is making message sending safe to retry, so a temporary network drop after you hit send does not create duplicate messages. One common way is to include a client generated message id that the server can treat as an idempotency key.
Sam: Do we really need both get messages and sync? Could we just poll get messages everywhere and skip the sync concept?
Dr. Wei: The read and sync bullet separates two needs: get messages is for a single thread, while sync is for reconciling multiple devices and catching up after gaps. The tradeoff is complexity versus efficiency: a dedicated sync path can be bandwidth-efficient and correct with cursors, but it needs careful state tracking. And the real time events bullet is about responsiveness: new message, typing, presence update, and receipt let us push changes without heavy polling, but it means we must manage long-lived connections and backpressure.
Data model and sharding strategy
Dr. Wei: Before we talk scaling, we need a clear data model and a sharding strategy that match the product experience. The goal is to support fast sends and reads, preserve message order within a conversation, and still make it easy to find what a given user should see in their inbox, even when message content is encrypted.
Sam: For sharding, why not shard by user id instead? That seems like it would keep a user inbox fast.
Dr. Wei: Now look at the conversation bullet and the partition bullet together: the conversation tracks participants and last sequence, and we shard by conversation id for ordering locality. The tradeoff is that sharding by conversation id makes writes and ordered reads in a thread simple and avoids cross-shard coordination for sequencing, but you still need an index by user id for inbox and notifications, which can be maintained asynchronously to avoid slowing the send path.
High-level architecture: online push vs offline queue
Dr. Wei: At a high level, this architecture separates two concerns: getting a message into the system reliably, and getting that message to the recipient quickly when they are online or safely when they are offline.
Sam: If we already publish to an event bus, can we skip writing to the message store first and just rebuild from the bus later?
Dr. Wei: The durability line answers that: the message store is the source of truth, and the event bus is a distribution mechanism that is typically at least once. In practice you connect them with an outbox or change data capture flow so publishing can be retried safely if it fails after the store write. Because the bus can redeliver, duplicates are expected and bounded, so downstream handlers and delivery consumers must be idempotent. Then the delivery line splits online push from offline queue and sync on reconnect, so online users get low latency pushes while offline users rely on queued delivery and cursor-based catch up.
Deep dive: message ordering (the thing users argue about)
Dr. Wei: Start with the rule bullet: strong ordering within a conversation id, and no global ordering. The tradeoff is we give users a consistent timeline where they care, inside one chat, while avoiding impossible coordination across all chats.
Sam: Could we just use timestamps from devices to order messages? That seems simpler than a leader per conversation.
Dr. Wei: That is the tempting wrong answer. The precondition bullet is there because device clocks drift and networks reorder; you need one sequencing authority per conversation shard. The implication is failover is the scary moment: if the leader changes, you bump an epoch or term and fence the old leader, otherwise you can accept conflicting sequences and users will literally see different histories.
Dr. Wei: Now the mechanism and multi-device bullets: the server assigns a monotonic sequence per conversation, so there should not be ties to break. For multiple devices and flaky networks, the key is idempotency: clients attach a stable client message id, and the server deduplicates on conversation id plus client message id, so retries and reconnect resends get the same assigned sequence instead of creating duplicates or reordering the timeline.
Deep dive: presence at 10 to the 9 scale
Dr. Wei: Start with the state bullet: presence is online, offline, and last seen per user id and often per device. The implication is you need a clear freshness policy, like how quickly you mark someone offline after missed heartbeats, because overly aggressive timeouts cause flapping and overly relaxed timeouts look stale.
Sam: Why not just broadcast every presence change to all of a user’s contacts? That would make the dots accurate.
Dr. Wei: That is exactly the anti-pattern bullet: broadcast to all contacts blows up like N times M, and it fails first on backend cost and battery, and it can also surprise users from a privacy standpoint. The fix bullet is to subscribe only a recent or relevant set, and do a lazy fetch when you open a conversation. The tradeoff is slightly less global freshness, but dramatically less fanout work, and it keeps presence updates scoped to where they are actually needed.
Delivery receipts $+$ offline sync
Dr. Wei: Start with the states bullet: sent to delivered to read. The implication is each state answers a different question, and if you blur them, you end up lying to users, especially during retries and offline periods.
Sam: If we already wrote the message durably, can we immediately mark it delivered? It is basically delivered to the system.
Dr. Wei: That confusion is why the acks bullet is explicit: server ack means sent to the service, device ack means it reached a recipient device, and read-event means the user opened it. Now add an important clarification for groups and multi-device: receipts are not a single bit on the message. They are keyed by message identifier and recipient identifier, and often also by device identifier. For a group, you either show per-recipient status, or aggregate it, like delivered to k of n recipients. You also have to define the minimum evidence to upgrade state, like delivered after any active device acks, or only after all of that user’s devices acks. The tradeoff is you need client cooperation and retry-safe bookkeeping. Then the offline bullet says we track a per-user pending cursor and use sync on reconnect. The implication is cursors let you resume efficiently and avoid re-sending an entire history, but you must be careful that cursors only move forward so state does not jump backward.
Group messaging: fan-out and ordering tradeoffs
Dr. Wei: Take the small groups bullet: eager fan-out to all members at send time. The implication is great read latency because everyone gets the message immediately, but the tradeoff is more writes per send and higher risk of partial fan-out during failures.
Sam: For large groups, should we still do eager fan-out if we have a big enough cluster? It feels simpler.
Dr. Wei: Now the large groups bullet: hybrid fan-out, online push plus pull on open. The tradeoff is more client work and slightly more read latency for offline members, but it caps backend work per send. And the ordering bullet says one sequence per group conversation id, with typing debounced. The implication is a single sequence keeps the group timeline consistent, while debouncing typing avoids turning a large group into a storm of tiny events.
Evaluate: failures, tradeoffs, and product evolution
Dr. Wei: To wrap up the design, we evaluate it through three lenses: the failures it must survive, the tradeoffs we make between consistency, latency, and availability, and how the system can evolve as the product adds new capabilities over time.
Sam: On tradeoffs, is it fair to say we always pick strong consistency if we are Meta-scale? Inconsistency seems unacceptable for chat.
Dr. Wei: Look at the tradeoffs bullet: push versus pull, and strong per conversation id versus eventual across conversation id. The tradeoff is you spend strong consistency budget where users notice it, inside a chat thread, and you relax it where it is less visible, across independent threads, to keep availability and scalability. Then the extensibility bullet: an event bus lets you add reactions, threads, and payments as new event types. The implication is versioned events and idempotent handlers let you evolve product features without rewriting the core message pipeline.
Exit ticket: sanity-check ordering and receipts
Dr. Wei: Look at the practice bullet: if the last seen sequence is one hundred four and the client retries with the same client message id X, you should think in terms of idempotency and monotonic progress. The implication is you can safely re-emit the server received acknowledgment for the same client message id, but you cannot claim delivered or read unless you have actual device or read evidence.
Sam: If the server already assigned a sequence earlier, should it reuse that same sequence on retry, or issue a new one and mark the old as superseded?
Dr. Wei: Reuse the same sequence for the same client message id, otherwise you create duplicates in the ordered log. Now the reflection bullet: if presence pushes to all contacts instead of a recent set, the first things to break are backend cost and battery, and you also expand the privacy surface. The design choice that limits blast radius is scoping subscriptions to a small relevant set and using lazy fetch when a conversation becomes active.
Thank you for watching!
Thanks for watching. Subscribe and share if you found this useful—see you next time!