Practice · Question 3 · Hard
Design a chat system (like WhatsApp)
Real-time one-to-one and group messaging with delivery receipts, offline delivery and multi-device sync. Tests WebSockets at scale, routing, ordering and delivery guarantees.
1. Clarify requirements
Functional
- 1:1 and group chats (groups up to ~500 members).
- Sent / delivered / read receipts.
- Offline delivery: messages wait and arrive when the user comes online; push notification meanwhile.
- Multi-device: history syncs across a user’s phone and laptop.
- Online presence and “last seen”.
- Out of scope for now: voice and video calls, stories. Media is mentioned briefly.
Non-functional
- Low latency delivery (< 200 ms when both users are online).
- No message loss, and ordering within a conversation.
- High availability; durable storage.
- Scale: 50M daily active users, ~40 messages per user per day.
2. Estimates
- Messages: 50M × 40 = 2 billion/day → 2 × 10⁹ ÷ 10⁵ ≈ 20k msg/s average, ~60k/s peak.
- Storage: 2B × ~200 bytes ≈ 400 GB/day → ~150 TB/year (before replication). A write-heavy, append-mostly workload.
- Concurrent connections: suppose 20% of DAU are online at peak → 10M open WebSockets. At ~50k connections per gateway server, that’s ~200 gateway servers.
3. API
Real-time events over the WebSocket:
→ send { client_msg_id, conversation_id, body, sent_at }
← ack { client_msg_id, message_id, seq } (stored)
← message { message_id, conversation_id, seq, sender, body }
→ receipt { conversation_id, up_to_seq, type: delivered|read }
← presence { user_id, online, last_seen }
REST for history and setup:
GET /v1/conversations/{id}/messages?before_seq=…&limit=50
GET /v1/sync?since_cursor=… (everything new for this device)
POST /v1/conversations (create group)
4. Data model
| Store | Data | Why |
|---|---|---|
| Message store (Cassandra / ScyllaDB) | messages(conversation_id PK, seq CK, message_id, sender_id, body, created_at) |
Write-heavy and append-only; reads are “latest N in a conversation”; partition by conversation, cluster by seq |
| SQL | users, conversations, members | Relational, low volume, needs constraints |
| Redis | session registry user → [gateway, device], presence, unread counts |
Fast, ephemeral, frequently updated |
| Per-user inbox | inbox(user_id, conversation_id, last_seq) |
What each user has seen; drives sync |
5. High-level design
6. Deep dives
The life of a message
- Alice’s app sends
sendwith aclient_msg_id(a UUID it generates) over her WebSocket. - The chat service assigns the next
seqfor the conversation, and persists the message. - Only after the write succeeds does it reply
ack. Alice’s app shows one tick (sent). - It looks up Bob in the session registry. If he’s online on gateway 2, it forwards the message there, and gateway 2 pushes it down Bob’s socket.
- Bob’s app replies
receipt delivered, and Alice sees two ticks. When Bob opens the chat, areadreceipt follows. - If Bob is offline, the message is already stored; the push service sends a notification, and Bob fetches it on reconnect.
Delivery guarantees and duplicates
Networks drop acknowledgements, so the client retries unacknowledged sends: at-least-once delivery. To avoid duplicates, the server keeps a short-lived record of client_msg_ids per conversation; a retry with a known ID returns the original ack instead of storing again. Recipients dedupe by message_id too.
Ordering
- A global clock order is impossible to guarantee, but per-conversation order is what users notice.
- Assign a monotonically increasing
seqper conversation, for example an atomic counter in the conversation’s partition, or by routing all writes for a conversation through one owner (consistent hashing onconversation_id). - Clients display by
seq, and can detect gaps (seq 41 then 43 means “fetch 42”).
Offline and multi-device sync
Each device stores the last seq it has seen per conversation (or a global cursor). On reconnect it calls GET /sync?since_cursor=… and gets everything newer. Delivered and read state is per device too. This is a simple, robust model: the message store is the source of truth and devices catch up from it.
Group chats: fan-out
- Write once, read by many: store each group message once in the conversation’s partition, then deliver to each online member’s gateway. For 500 members that’s up to 500 gateway pushes. Do it asynchronously through a queue, so the sender gets an acknowledgement immediately.
- Look up members’ sessions in batches; push to offline members only as notifications, rate-limited and collapsed (“12 new messages”).
- Very large channels (100k members) work better pull-based: clients fetch on open instead of receiving every message pushed.
Presence
- Online status comes from the WebSocket itself plus heartbeats (say every 30 s); missing heartbeats mean offline after a grace period, which avoids flicker on bad networks.
- Presence fan-out is expensive (each change × every contact). Mitigate: only send presence to users who currently have the chat open (subscribe on view), and batch updates.
Scaling the gateway tier
- Gateways are stateful (they hold sockets) but do no business logic, so they’re cheap to add.
- Use least-connections load balancing; on deploys, drain gateways gradually so clients reconnect smoothly with backoff and jitter (avoids a reconnect stampede).
- The session registry updates on connect and disconnect, with TTLs so crashed gateways’ entries expire.
Media and security
- Images and video: upload to object storage with pre-signed URLs; the message carries only the object key and a thumbnail.
- End-to-end encryption (the Signal protocol): the server routes and stores ciphertext it can’t read. That affects features like server-side search, which is worth mentioning as a trade-off.
7. Bottlenecks and failure modes
| What fails | Impact | Mitigation |
|---|---|---|
| A gateway dies | ~50k users disconnect | Clients reconnect to another gateway (backoff + jitter) and sync from their cursor |
| Session registry stale | Messages routed to the wrong gateway | TTLs + heartbeats; gateway rejects unknown users → fall back to push |
| Message store partition slow | Sends delay | Replication; tune consistency (QUORUM writes); backpressure to clients |
| Push provider outage | Offline users not notified | Retry with backoff; messages are safe in storage and sync on next open |
| Huge group spam | Fan-out storm | Async fan-out queue, per-group rate limits, pull model for big channels |
8. Wrap-up
A strong answer: a gateway tier for WebSockets plus a session registry; persist then acknowledge; per-conversation sequence numbers; client message IDs for idempotent retries; cursor-based multi-device sync; a deliberate group fan-out strategy; and push notifications for offline users.
Likely follow-ups: How would you add message search (and how does E2E encryption limit it)? How do you delete a message for everyone? How would you support 100k-member channels? How do you keep the gateway tier from melting during a regional reconnect storm?
Test yourself
Answer in your head, then click a card to check. All cards are in the Anki deck.