Course Content
System Design Interview
31 sections · 71 lessons
Chat System: message storage, groups, presence and follow-ups
Chat System: scope, scale, connections and the stateful tier got a message from the sender's server to the recipient's. This lesson covers what happens to it there: where it is stored and in what order, how it reaches a group, the presence and delivery state around it, and the extensions interviewers ask about most.
Storage comes down to two decisions: what database holds 600 GB a day forever, and what makes messages appear in the same order for everyone in a conversation.
Message storage and ordering
Why not a relational database
The access pattern is narrow and the volume is large:
- Write rate: 60,000 inserts/s at peak, append-only, never updated except for state flags.
- Read pattern: "give me the last 50 messages in conversation C before message ID X." A range scan within one partition. No joins, no aggregation, no ad-hoc queries.
- Volume: 220 TB/year and growing without bound.
A relational database can be made to do this with heavy sharding, but you are then paying for transactions, joins, and secondary indexes you never use, and doing the partitioning yourself. A wide-column store — the family that includes Cassandra and HBase — is built for exactly this shape: high write throughput, a partition key plus a clustering key, and range scans inside a partition (Databases: choosing and justifying).
The data model
messages partition key: conversation_id clustering key: message_id (descending) columns: sender_id, body, created_at, attachment_ref, statePartitioning by conversation puts every message in a conversation on one node, which makes the range scan a single-node read. It also creates the one hazard: a conversation with millions of messages becomes an unbounded partition. The standard fix is a composite partition key such as (conversation_id, time_bucket) where the bucket is a month, capping partition size and costing an extra lookup when a read crosses a bucket boundary.
Ordering, and why timestamps fail
The obvious ordering key is the server's clock. It does not work, for two reasons that compound.
Clock skew. Two chat servers in the same datacentre, both synchronised by the Network Time Protocol, can still differ by single-digit to tens of milliseconds. Messages sent 5 ms apart through two different servers can be stamped in the wrong order — and the receiver sees a reply appear above the message it answers.
Collisions. At 60,000 messages per second, millisecond timestamps collide constantly. Two messages with the same timestamp have no defined order, so two clients sorting the same data can produce two different conversations.
What to use instead
Order within a conversation is the only order that has to be right. Nobody can observe the relative order of two messages in two different conversations.
- A per-conversation sequence number. The conversation's owning partition assigns 1, 2, 3. Perfectly ordered and gap-free, at the cost of a coordination point per conversation. Viable because conversations are independent and mostly low-rate.
- A globally unique, time-sortable identifier of the Snowflake family from The Snowflake approach: a timestamp in the high bits, a machine identifier, and a per-machine sequence in the low bits. Sortable, no coordination, and unique even under clock collision.
Recommend the time-sortable identifier as the default — no per-conversation coordination point, and it doubles as the pagination cursor and the client's sync marker. Use per-conversation sequence numbers when the product needs gap detection ("am I missing message 41?"), which the delivery state below can exploit.
Group chat, presence, and delivery state
Three features that look like product details and each contain a real scaling decision.
Group chat: fan-out, again
Section 13 established the fan-out on write versus on read trade-off. Chat resolves it differently, and it is worth saying why rather than repeating the analysis.
For a group of 100, fan-out on write wins outright: one message becomes 100 inserts into 100 per-user inboxes. The recipient's read path is then a single scan of their own inbox, merged across all conversations, which is exactly what a chat client needs — one ordered stream, not 40 separate conversation queries.
The threshold where this fails is much higher than in a news feed, because groups have caps. At the documented limits of typical messaging products — hundreds to a few thousand members — 100 to 1,000 inserts per message is affordable at 60,000 messages/s only if large groups are rare, which they are. For broadcast channels with a million subscribers, switch to fan-out on read exactly as in The hybrid that real systems use.
Presence, and the flapping problem
Presence is "is this person online?" The naive implementation is the connection itself: online means a socket exists. Two problems.
Detecting the difference between "closed the app" and "went into a tunnel." A clean disconnect sends a close frame. A phone losing signal sends nothing. The fix is a heartbeat: the client sends a small ping every 10 seconds, and presence expires 30 seconds after the last one. The cost is 15M ÷ 10 = 1.5 million heartbeats per second, which is 25× the message rate. Presence is more traffic than messaging.
Flapping. A user on a train oscillates between connected and disconnected every few seconds. Each transition, published naively, is a fan-out to everyone who can see their presence. A user with 500 contacts flapping ten times a minute generates 5,000 presence updates a minute from one phone.
Three fixes, applied together: a grace period before publishing "offline" (30 to 60 seconds), debouncing so at most one update per user per interval is published, and publishing presence only to contacts who currently have that user's conversation open — which cuts the fan-out audience by an order of magnitude.
Delivery state
The sent → delivered → read progression is a small state machine with one message per transition:
| State | Set when | Who reports it |
|---|---|---|
| Sent | The server persisted the message | Server acknowledges to the sender |
| Delivered | The recipient's device received it | Recipient's client sends an acknowledgement |
| Read | The recipient opened the conversation | Recipient's client sends a read receipt |
Each transition is itself a message that has to reach the original sender, so a conversation generates roughly three times its message count in traffic. In a group of 100, read receipts from every member multiply that by 100 — which is why group read receipts are often reduced to a count, or dropped entirely above a group size. Batch acknowledgements ("I have everything up to message X") rather than sending one per message.
Follow-ups
The four extensions that come up most, each with the honest answer.
Offline delivery and the push handover
A message arrives for a user with no connection. The message is already durable in the store and referenced in that user's sync queue, so nothing is lost. What remains is waking them up.
The chat service calls the notification system from Section 12, which delivers through Apple Push Notification service or Firebase Cloud Messaging. Three details that matter:
- The push carries a hint, not the message. A notification payload is size-limited (a few kilobytes) and passes through a third party. Send "you have a new message from Ravi" and let the client fetch the content on open. With end-to-end encryption, the server could not include the content anyway.
- Deduplicate the wake-up. Forty messages arriving in a minute should be one notification, which is exactly the digest window from User preferences and throttling.
- The reconnect is the real delivery. The client reconnects, sends its last known message identifier, and the server replays the sync queue from there. Push is a doorbell, not a delivery mechanism.
Media attachments
Media never travels through the message path. The client uploads to blob storage using a pre-signed URL, receives an identifier, and sends a message whose body is that identifier plus a thumbnail. Recipients fetch through a content delivery network.
The reason is arithmetic: a 2 MB photo through the message pipeline at even 1% of message volume would be 600 messages/s × 2 MB = 1.2 GB/s through servers designed to move 300-byte rows. The upload path, in Section 16, develops this properly.
End-to-end encryption and what it takes away
With end-to-end encryption, keys live on devices and the server stores ciphertext. That is a strong privacy property with four concrete architectural consequences:
- No server-side search. Search must run on the device, over locally stored messages, which means the client keeps a full local index and search quality is bounded by what fits on a phone.
- No server-side moderation of content. Abuse handling shifts to metadata signals and user reports.
- Multi-device sync becomes hard. Each device has its own keys, so a message must be encrypted once per recipient device, and adding a device requires a key-transfer protocol.
- Server-side previews and notifications lose content. The push notification cannot say what was said unless the client decrypts and rewrites it locally.
None of these is a reason not to do it. They are the price, and naming them precisely is what a senior answer looks like.
Multi-device sync
One account, four devices. The sync queue from the group chat section above becomes per device, not per user: each device tracks its own last-received identifier, because a phone that was off for a week and a laptop that was online five minutes ago need different amounts of replay.
Read state, in contrast, is per user — marking a conversation read on the laptop must clear the badge on the phone. So the model splits: message delivery is per device, conversation state is per account, and both propagate over the same connection.