System design · Queues, write scaling, delivery lag

How to design a chat system like WhatsApp

A chat system looks simple: a message goes from one person to another. The hard part is doing it for millions of people at once - holding open connections, keeping every conversation in order, delivering to people who are offline, and surviving the moment everyone sends "Happy New Year" in the same second.

Updated · 6 min read

Requirements

Agree the scope first. Chat has many optional features; interviewers want to see you pick the core and say what you're leaving out.

  • Functional: one-to-one messages; group chats; message history on every device; sent, delivered and read receipts.
  • Functional, optional: online presence and "last seen"; typing indicators; media attachments; push notifications for offline users.
  • Non-functional: a sent message should be acknowledged in well under a second and delivered within seconds; messages in a conversation must arrive in order and never be lost or duplicated.
  • Non-functional: the system must absorb sharp spikes - holidays, live events - without dropping messages.

Capacity estimates

These are assumptions, stated out loud. The point is the order of magnitude, which decides the architecture.

QuantityAssumptionResult
Messages100M daily users × 40 messages each4B a day ≈ 46,000/s average, ≈ 140,000/s at a 3× peak
History readsabout 1.5 reads per message sent≈ 70,000/s average
Storage≈ 200 bytes per message, text only≈ 800 GB a day ≈ 290 TB a year, before replicas
Open connections30% of daily users online at peak≈ 30M concurrent connections

Three conclusions. Writes are a large share of the load - unlike a URL shortener, caching alone won't save you. Storage grows by hundreds of terabytes a year, so the message store must be partitioned from day one. And tens of millions of long-lived connections need their own tier of servers.

API design

Clients keep a WebSocket open for real-time traffic in both directions, and use plain HTTP for history, which is cacheable and easy to page through.

WebSocket /ws
→ { "type": "send", "client_msg_id": "c-91f2", "conversation_id": "conv_42", "body": "hi" }
← { "type": "ack", "client_msg_id": "c-91f2", "message_id": "1840238472101", "status": "sent" }
← { "type": "message", "conversation_id": "conv_42", "message_id": "1840238472101", "from": "u_7", "body": "hi" }
→ { "type": "receipt", "message_id": "1840238472101", "status": "read" }

GET /api/conversations/{id}/messages?before={message_id}&limit=50
→ 200 OK  { "messages": [ ... ] }

The client-generated client_msg_id matters: if the connection drops before the ack arrives, the client resends with the same ID and the server can drop the duplicate.

Data model

The dominant query is "the latest messages in this conversation", so messages are partitioned by conversation and sorted by a time-ordered message ID. That shape fits a wide-column store such as Cassandra, or a key-value store such as DynamoDB, or PostgreSQL sharded by conversation.

messages
  conversation_id  bigint     partition key
  message_id       bigint     sort key (time-ordered, e.g. Snowflake)
  sender_id        bigint
  body             text
  created_at       timestamp

conversation_members
  conversation_id  bigint
  user_id          bigint
  last_read_id     bigint     for read receipts and unread counts

Never store receipts as updates to the message row: a group message read by 200 people would become 200 writes to one hot row. Store a per-member last_read_id instead.

High-level design

Split the system by job, and keep the slow parts off the path the sender waits on.

  • Connection gateways hold the WebSockets. They're stateful, so a routing store (often Redis) records which gateway each online user is connected to.
  • The chat service validates a message, assigns a message ID, and puts it on a queue partitioned by conversation - then acknowledges the sender. The sender waits only for this.
  • Workers consume the queue: they write the message to the message store and look up each recipient's gateway to deliver it. Recipients who are offline get a push notification instead.
  • History reads go through a separate read path with a cache in front of the message store, because recent messages are read far more often than old ones.

Where it breaks

Write every message synchronously to the database before acknowledging it, and the database becomes the bottleneck the moment traffic spikes. Writes can't be cached, and read replicas don't add write capacity, so at midnight on New Year's Eve the sends simply fail.

The fix is to decouple: accept the message onto a durable queue, acknowledge, and let workers write at a rate the store can sustain. A spike becomes a backlog - delivery is a few seconds late - instead of errors. The new risk is lag: if workers can't drain the backlog, delivery falls further behind until it's useless. Size workers so the backlog clears soon after the peak, and shard the message store so the workers' writes fit.

Message ordering and exactly-once delivery

  • Order per conversation, not globally. Partition the queue by conversation ID so one worker handles a conversation's messages in sequence.
  • Use time-ordered message IDs (a Snowflake-style generator) so storage and clients sort the same way.
  • Deduplicate with the client's client_msg_id. Networks retry; idempotency turns "at least once" into "exactly once" from the user's point of view.

Delivery receipts and offline users

  • Sent: the server acknowledged and queued the message. Delivered: a recipient's device received it. Read: the recipient opened the conversation.
  • An offline device syncs on reconnect by asking for everything after its last known message ID in each conversation.
  • Push notifications wake offline phones; they carry a hint, not the source of truth - the app still syncs from the server.

Group chats and fan-out

A message to a group of 200 means 200 deliveries. For small groups, fan out on write: workers deliver to every online member. For very large channels, fan out on read: store the message once and let members pull it when they open the channel, so one message doesn't become millions of deliveries.

Presence

Online status is the chattiest feature: every connect, disconnect and heartbeat is an update. Store presence with a short expiry, refreshed by heartbeats, and only publish changes to people who are looking - a user's contacts with the app open - rather than to everyone who has them in a list.

What interviewers look for

  • Separating the connection tier from the logic and storage tiers, and explaining why gateways are stateful.
  • Acknowledging after a durable queue, not after the database write - and knowing what that costs (delivery lag) and how to bound it.
  • Per-conversation ordering, idempotent sends, and a storage design partitioned by conversation.
  • A clear answer for group fan-out and offline delivery, with numbers.

Frequently asked questions

Why use WebSockets for chat instead of polling?

+

Chat needs the server to push messages the moment they arrive. Polling either wastes requests when nothing has changed or adds delay between polls. A WebSocket keeps one connection open in both directions, so delivery is immediate and idle users cost almost nothing.

How do you keep messages in order?

+

Order is only needed within a conversation. Route every message of a conversation to the same queue partition so one consumer processes them in sequence, and give each message a time-ordered ID so clients and storage sort identically.

How do you store billions of messages?

+

Partition by conversation ID and sort by message ID inside each partition. Recent messages are read far more than old ones, so cache recent history and move old messages to cheaper storage over time.

How is group chat different from one-to-one chat?

+

Fan-out. A one-to-one message is one delivery; a group message is one per member. Small groups fan out on write; very large channels store the message once and let members fetch it on read.

How does the server know where to deliver a message?

+

Each connection gateway registers which users are connected to it in a shared routing store. A worker delivering a message looks up the recipient's gateway and hands the message to it; if the user isn't connected anywhere, it sends a push notification instead.

Now break one yourself.

The first challenge takes about two minutes. No signup.