Limited time: AI code review, hints, mock interviews, whiteboard analysis, and all Pro features are unlocked. Enroll
โฑ๏ธ 40 min read

Designing a Chat System Like WhatsApp / iMessage

Difficulty: Intermediate Prerequisites:Message Queues, Caching, and WebSockets


TL;DR

A chat system delivers messages in real-time using WebSockets for online users and stores-then-forwards for offline users.

flowchart LR
    SENDER["Sender"]:::client
    WS["WebSocket Servers"]:::service
    CHAT["Chat Service"]:::service
    STORE[("Message Store<br/>Cassandra")]:::data
    K["Kafka<br/>fan-out"]:::async
    PUSH["Push Notifications<br/>FCM APNs"]:::external
    RECEIVER["Receiver"]:::client

    SENDER -->|"1. Open WebSocket"| WS
    WS -->|"2. Forward message"| CHAT
    CHAT -->|"3. Persist message"| STORE
    CHAT -->|"4. Publish to fan-out"| K
    K -->|"5. Fan out"| WS
    CHAT -->|"6. Push offline alert"| PUSH
    WS -->|"7. Deliver to recipient"| RECEIVER

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef async fill:#AB47BC,stroke:#4A148C,color:#fff
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
    classDef external fill:#4a1942,stroke:#f472b6,color:#e2e8f0

In 3 sentences: Clients maintain a persistent WebSocket connection to the server. When a message is sent, the server persists it, looks up which server the receiver is connected to, and pushes it down their WebSocket. If the receiver is offline, the message waits in a queue and a push notification is sent.


Understanding the Problem

๐Ÿ’ฌ What is a chat system? A real-time messaging platform that lets users send text, images, and files to individuals or groups. Messages must be delivered reliably (even if the recipient is offline), ordered correctly, and displayed in real-time. Think WhatsApp, Telegram, Facebook Messenger, or Slack. The hard parts: guaranteed delivery across flaky mobile networks, real-time push without polling, group fan-out at scale, and end-to-end encryption.

Naive First Cut

flowchart LR
    SENDER["Sender"]:::client
    API["Chat API"]:::service
    DB[("Messages DB<br/>one table")]:::data
    RECEIVER["Receiver<br/>polls every 5s"]:::client

    SENDER --> API
    API --> DB
    RECEIVER --> API

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0

Sender POSTs message to an API, stored in a DB. Receiver polls the API every 5 seconds for new messages.

Why this breaks:

Prior Art Weโ€™re Drawing From

Functional Requirements

Core:

  1. Users can send messages (text) to another user in real-time (1:1 chat).
  2. Users can create groups and send messages to all group members.
  3. Messages are delivered reliably even if the recipient is offline (store-and-forward).

Below the line:

Non-Functional Requirements

Core:

Below the line:

Scale Estimation (Back-of-Envelope)

Three numbers drive everything below: 50M concurrent connections to hold open, 1.2M messages/sec to route, and a P99 delivery budget of 500ms.

Core Entities

API / System Interface

WebSocket: wss://chat.example.com/ws
  โ†’ Client authenticates on connect (JWT)
  โ†’ Bidirectional: send messages, receive messages, typing, presence

REST (fallback + media):
POST /v1/messages         โ†’ send a message (fallback if WS down)
GET  /v1/conversations/:id/messages?after=<seqNo>  โ†’ sync history
POST /v1/media/upload     โ†’ upload image/file, get a mediaUrl
POST /v1/groups           โ†’ create group

Wire format (over WebSocket). Worth splitting the two directions, because who may state an identity differs between them.

Client โ†’ server. Note there is no sender field anywhere in here:

{"type": "message", "to": "conv_123", "text": "hello", "clientMsgId": "uuid"}
{"type": "typing", "conversationId": "conv_123"}
{"type": "read", "conversationId": "conv_123", "upToSeqNo": 4821}

Server โ†’ client. The server attaches identity, having derived it from the connectionโ€™s authenticated session:

{"type": "message", "messageId": "msg_456", "conversationId": "conv_123", "senderId": "u_789", "text": "hello", "seqNo": 4822}
{"type": "ack", "clientMsgId": "uuid", "messageId": "msg_456", "status": "delivered"}
{"type": "typing", "conversationId": "conv_123", "userId": "u_789"}

Security: the WebSocket is authenticated via JWT on handshake, and the server holds the resulting identity against the connection for its lifetime. senderId is taken from that session, never from the frame โ€” otherwise any client can send as anyone. That same rule is why the inbound typing frame carries no userId: a client able to name the typist could forge โ€œAlice is typingโ€ in any conversation it can reach. The server stamps userId on the way out. It also generates the authoritative messageId and seqNo; clientMsgId is only the senderโ€™s own reference, echoed back in the ack so the client can reconcile its optimistic render.

High-Level Design

1) User sends a 1:1 message (both online)

Build the simplest thing that delivers a message from Alice to Bob in real time. One server and one database. Everything else on this page โ€” the connection registry, the sequence numbers, the offline queue, the event bus โ€” answers a non-functional requirement, so each one belongs in a deep dive where we can show why it is needed instead of asserting it up front.

New components we need:

  1. Chat Server - holds an open WebSocket to every online user and does the routing. Because there is exactly one of these, it can keep the whole userId โ†’ connection table in a plain in-memory map: if Bob is online, this process is holding his socket.
    ๐Ÿ’ก WebSocket = a persistent connection that stays open so the server can push messages instantly without the client asking. Unlike HTTP (ask โ†’ answer โ†’ done), WebSocket keeps the line open. Learn more โ†’
  2. Message Store (Cassandra) - permanent storage for every message, partitioned by conversationId so loading a chat history is a single-partition read.
flowchart LR
    SENDER["Sender"]:::client
    CHAT["Chat Server<br/>holds all sockets"]:::service
    MAP["In-memory<br/>userId to socket"]:::data
    STORE[("Message Store<br/>Cassandra")]:::data
    RECEIVER["Receiver"]:::client

    SENDER -->|"1. Send over WebSocket"| CHAT
    CHAT -->|"2. Persist message"| STORE
    CHAT -->|"3. Look up receiver socket"| MAP
    CHAT -->|"4. Push down that socket"| RECEIVER
    RECEIVER -->|"5. Delivered ack"| CHAT
    CHAT -->|"6. Relay ack to sender"| SENDER

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0

Step-by-step flow:

  1. Sender types โ€œHey, are you free tonight?โ€ and hits send โ†’ the frame travels over their already-open WebSocket to the Chat Server
  2. Chat Server writes the message to the Message Store with partition key conversationId. Store before delivering โ€” once this write returns, the message survives any crash, and that is what lets us ack the sender honestly
  3. Chat Server looks up the receiver in its in-memory connection map โ†’ found, Bob is connected
  4. It writes the message frame down Bobโ€™s socket. It appears on his screen
  5. Bobโ€™s app replies with a delivered ack
  6. Chat Server relays that back to Alice, who now sees the double-check โœ“โœ“

Why WebSocket instead of HTTP polling? This one genuinely is a functional requirement, not an optimization. Polling every 2 seconds for 500M users is 250M requests/second of mostly-empty responses, and it puts a floor of 2 seconds on delivery latency against a 500ms P99 target. There is no version of โ€œreal-timeโ€ that polling reaches, so the persistent connection is in the base design rather than deferred.

What we have deliberately left broken. For a chat app with a few thousand users on one box, the design above is complete and correct. It has three holes:

2) Receiver is offline - store and forward

Here is the useful thing about having stored the message first: when Bob is offline, most of the work is already done. The message is durable in the Message Store whether or not anyone is listening. What is left is a way to wake Bobโ€™s phone, and a way for Bob to find out what he missed when he comes back.

New components we need (in addition to the ones above):

  1. Push Service - sends push notifications to wake up the userโ€™s phone.
    ๐Ÿ’ก Think of it as the โ€œtap on the shoulderโ€ that tells the user to open the app.
  2. FCM / APNs - Firebase Cloud Messaging (Android) and Apple Push Notification service (iOS). External services that deliver notifications to locked phones.
    ๐Ÿ’ก FCM doesnโ€™t โ€œknowโ€ a message arrived - YOUR server tells FCM to send the push. When Bob installs the app, FCM gives his device a unique token. Your server stores this token. When Bob is offline and a message arrives, your server calls FCMโ€™s API with Bobโ€™s token + notification content. FCM maintains its own persistent connection to every Android device in the world and routes the push through that always-on channel. APNs works the same way for iOS. Learn more about real-time communication โ†’

No offline queue. We do not need a second copy of the message somewhere else, because the Message Store is already partitioned by conversation โ€” โ€œwhat did I miss in this chatโ€ is a read of that partition.

How does the notification show the actual message text (with E2E encryption)?

For E2E encrypted apps like WhatsApp/Signal, FCM does NOT carry the message content (the server canโ€™t read it). Instead:

  1. Server sends a silent data message via FCM - just a โ€œwake up, you have a new messageโ€ signal with sender ID and message reference
  2. FCM wakes up the appโ€™s background process on the device
  3. The app connects to the server, pulls the encrypted message, and decrypts it locally on the device
  4. The app constructs the notification itself (โ€œAlice: Hey, are you free?โ€) and hands it to the OS for display

For non-E2E apps, the server CAN send the message text directly in the FCM payload (notification message type) - simpler but less secure.

flowchart LR
    SENDER["Sender"]:::client
    CHAT["Chat Server"]:::service
    STORE[("Message Store")]:::data
    PUSH["Push Service"]:::service
    FCM["FCM and APNs"]:::external
    RECEIVER["Receiver later"]:::client

    SENDER -->|"1. Send message"| CHAT
    CHAT -->|"2. Persist message"| STORE
    CHAT -->|"3. Receiver not in socket map"| PUSH
    PUSH -->|"4. Wake the device"| FCM
    FCM -->|"5. Notification"| RECEIVER
    RECEIVER -->|"6. Reconnect and ask what changed"| CHAT
    CHAT -->|"7. Read conversation partition"| STORE

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
    classDef external fill:#4a1942,stroke:#f472b6,color:#e2e8f0

Step-by-step flow:

  1. Alice sends the message exactly as in flow 1
  2. Chat Server persists it to the Message Store first, as always โ€” store, then deliver
  3. It looks up Bob in the connection map and does not find him. Bob is offline
  4. Push Service asks FCM or APNs to wake Bobโ€™s device
  5. Bobโ€™s phone buzzes with โ€œNew message from Aliceโ€
  6. Whenever Bob next opens the app, his client reconnects and asks for everything it does not already have in each conversation
  7. Chat Server reads the conversation partition and streams back the messages Bob is missing

Why store-and-forward instead of just โ€œretry laterโ€? Mobile networks are unreliable. A user might be offline for hours - on a flight, in a tunnel, phone dead. Because step 2 happens before any delivery attempt, the guarantee is unconditional: once the server acks Alice, the message will reach Bob eventually, however long that takes. Retrying in memory would lose it on the next deploy.

What we have deliberately left broken. Step 6 is doing a lot of hand-waving. โ€œEverything it does not already haveโ€ needs a cursor, and this design has nothing to put in one โ€” no ordering, so no well-defined โ€œafter this point.โ€ Today the client would have to fall back to โ€œsend me the last N messages and I will diff them locally,โ€ which is wasteful and gets worse the longer Bob was away. Deep Dive 2 gives us the sequence number that makes a cursor possible, and Deep Dive 4 turns it into a real sync protocol. That is also where a Redis-backed offline queue starts to earn its place, as a fast index of which conversations changed while Bob was gone, so reconnect does not have to poll all of them.

3) Group message fan-out

A group is a conversation with more than two participants. That is genuinely all it is, and the storage side needs nothing new: the message already lives in one partition keyed by conversationId, and a group conversation has one of those the same as a 1:1 chat does.

New components we need: none. What we need is one more table.

  1. conversation_members - which users belong to which conversation. This is the only new thing a group requires, and everything else reuses flow 1 and flow 2 per member.
flowchart LR
    SENDER["Sender"]:::client
    CHAT["Chat Server"]:::service
    STORE[("Message Store<br/>one copy")]:::data
    MEM[("conversation_members")]:::data
    MEMBERS["Online members"]:::client
    PUSH["Push Service"]:::service

    SENDER -->|"1. Send group message"| CHAT
    CHAT -->|"2. Store single copy"| STORE
    CHAT -->|"3. Read member list"| MEM
    CHAT -->|"4. Push to each online member"| MEMBERS
    CHAT -->|"5. Wake each offline member"| PUSH

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0

Step-by-step flow:

  1. Sender sends a message to group conv_123, which has 256 members
  2. Chat Server stores ONE copy with partition key conv_123, not 256 copies
  3. It reads the 256 member IDs from conversation_members
  4. It loops over them. For each member found in the connection map, write the frame down their socket
  5. For each member not found, call Push Service โ€” flow 2, once per absent member
  6. Ack the sender once the loop finishes

Why store once rather than a copy per member? 256 copies of the same text is 256x the storage for no read benefit, since a group chatโ€™s history is exactly the thing that is shared between all of them. It also makes edit and delete one write instead of 256.

What we have deliberately left broken. Step 4 and 5 run inside the senderโ€™s request, and that is where this falls apart:

Both are non-functional failures of a functionally correct design, and both are Deep Dive 5.

Technology Choices

Tier Purpose Primary pick Alternatives
Real-time transport Push messages to online clients WebSocket (long-lived) SSE, MQTT (IoT/mobile-optimized), gRPC streaming
Connection management Track whoโ€™s online on which server Redis Pub/Sub + connection registry Kafka, custom session store
Message storage Durable, ordered message log Cassandra (partition per conversation) ScyllaDB, DynamoDB, TiDB
Message queue Decouple sender from fan-out Kafka (per-user topic or partition) SQS, RabbitMQ, Pulsar
Offline delivery Store messages until recipient connects Redis sorted set per user SQS per user, Cassandra unread table
Presence Whoโ€™s online Redis with TTL per user Dedicated presence service
Media storage Images, files, voice notes S3 / GCS with CDN MinIO, Azure Blob
Push notifications Offline users FCM + APNs OneSignal, SNS
E2E encryption Message privacy Signal Protocol (Double Ratchet) Custom, Noise Protocol

Data Modeling

Cassandra (Message Store โ€” durable, ordered message log):

Table: messages
  PK: conversation_id
  SK: sequence_number (ascending)
  Columns: message_id (UUID), sender_id, content, type (text/image/file), 
           media_url, status (sent/delivered/read), created_at, client_msg_id

Table: conversations
  PK: user_id
  SK: last_message_at (DESC)
  Columns: conversation_id, participants (set), type (1:1/group), 
           unread_count, last_message_preview

Redis (Connection Registry โ€” whoโ€™s online where):

Key: "conn:{userId}" โ†’ Hash { serverId, connectionId, connectedAt }
TTL: refreshed on every heartbeat (30s), expires if client disconnects

Redis (Offline Queue โ€” pending messages for offline users):

Key: "offline:{userId}" โ†’ Sorted Set (score = seqNo, member = messageId)
Drained on reconnect, removed after delivery ack

Redis (Sequence Counter โ€” message ordering):

Key: "seq:{conversationId}" โ†’ Integer (atomically incremented via INCR)

Access Patterns:

Query Data Source How
Send message (persist) Cassandra INSERT INTO messages with assigned seqNo
Load chat history Cassandra Range query on messages where conversation_id = ? AND seqNo > lastSynced
Find receiverโ€™s server Redis HGETALL conn:{userId} โ†’ get serverId
Drain offline messages Redis ZRANGEBYSCORE offline:{userId} 0 +inf โ†’ deliver all, then DEL
List conversations Cassandra Query conversations by user_id ordered by last_message_at
Get next seqNo Redis INCR seq:{conversationId} โ†’ atomic, no contention across conversations

How Message Ordering Is Guaranteed in a Group Chat:

  1. Sender sends message โ†’ Chat Service calls INCR seq:{conversationId} โ†’ gets seqNo 42
  2. Message stored in Cassandra with (conversation_id, seqNo=42) as the clustering key
  3. Fan-out workers deliver messageId + seqNo to each online member via WebSocket
  4. Receiverโ€™s client inserts message at correct position by seqNo (not arrival time)
  5. If client receives seqNo 44 but hasnโ€™t seen 43 โ†’ it knows thereโ€™s a gap, requests the missing message: GET /conversations/:id/messages?after=42&before=44
  6. All devices sync via the same seqNo cursor โ€” each device tracks lastSyncedSeqNo and pulls the delta on reconnect

Deep Dives

1) Where do 50M persistent connections live, and how does a sender find the right one?

Problem: Flow 1 put every WebSocket on one Chat Server and kept the userId โ†’ socket table in that processโ€™s memory. Two things break, for different reasons, and it is worth separating them.

Capacity. 50M concurrent connections do not fit on one machine under any thread model.

Correctness. The in-memory map is only right because there is one server. Add a second and Aliceโ€™s server has no way to reach a socket held by Bobโ€™s server, so step 3 of flow 1 simply fails to find people who are demonstrably online. Scaling out is not an optimization here โ€” it breaks routing, and something has to put it back.

Bad โ€” the single server from flow 1, on the obvious thread model.

Take flow 1 as written and give each connection its own thread, which is what a straightforward blocking-I/O server does. A JVM thread stack is ~512KB, so 50M connections would need ~25TB of RAM in stacks alone, and you run into OS thread and file-descriptor limits four orders of magnitude before you get near that. Even the 2M connections we would like from one large node comes to ~1TB. This is not a tuning problem, the model is wrong: 50M connections are almost all idle almost all the time, and a thread is an expensive way to represent waiting.

The single server has a second failure that no thread model fixes. It holds 100% of live connections, so one deploy or one crash disconnects every user at once โ€” and destroys the routing table in the same instant, because the map lives in the same process. Then 50M clients reconnect simultaneously.

Good - NIO event loop model (Netty, Node.js, Go goroutines), across a fleet with a shared connection registry.

Instead of one thread per connection, use a small pool of threads (event loops) that multiplex thousands of connections using OS-level I/O selectors (epoll on Linux, kqueue on macOS). That fixes capacity. To fix routing, move the connection map out of process into a shared Connection Registry so any server can find any socket. This is the piece flow 1 could do without and a multi-server design cannot.

What Netty is: An asynchronous, event-driven network framework for Java. It implements the Reactor pattern - a single thread monitors many sockets, and only wakes up when thereโ€™s data to read/write. No blocking, no idle threads.

Event Loop (1 thread) monitors 100K connections via epoll
    โ†’ Connection has data? โ†’ Read it, process, respond
    โ†’ Connection idle? โ†’ Costs nothing (just a file descriptor)

Real numbers:

Tech used in production:

Great - tiered architecture separating connection from logic.

At extreme scale (100M+ connections), even Netty hits limits on a single machine. The problem: the server holding connections ALSO processes messages (routing, persistence, fan-out). Under load, message processing slows down AND connection handling suffers - they compete for the same CPU/memory.

The solution: split into two independent layers, each doing one job.

Edge Tier (connection holding) - the โ€œreceptionistโ€:

Logic Tier (message processing) - the โ€œbrainโ€:

How a message flows through both tiers:

  1. Aliceโ€™s phone is connected to Edge Server #3 via WebSocket
  2. Alice sends โ€œHey Bobโ€ โ†’ Edge Server #3 receives the raw bytes
  3. Edge Server #3 forwards to Logic Tier via internal gRPC: โ€œmessage from userId=alice, payload=Hey Bobโ€
  4. Logic Tier stores in Cassandra, then checks Connection Registry: โ€œBob is on Edge Server #7โ€
  5. Logic Tier sends to Edge Server #7: โ€œdeliver this to Bobโ€™s WebSocket connectionโ€
  6. Edge Server #7 pushes the message down Bobโ€™s WebSocket
  7. If Bob is offline โ†’ Logic Tier calls Push Service instead (FCM/APNs)

The Connection Registry (Redis) ties both tiers together:

Redis Hash: connection_registry
  alice โ†’ edge-server-3:conn-8842
  bob   โ†’ edge-server-7:conn-1204
  carol โ†’ edge-server-3:conn-9921

When Logic Tier needs to deliver to Bob, it looks up this registry and routes to the correct edge server. When Bob disconnects, Edge Server #7 removes the entry.

Why this is better than one server doing everything:

Real-world implementations:


2) Alice sends two messages a second apart. How do we guarantee Bob sees them in that order?

Problem: Nothing in the high-level design assigns an order. Flow 1 renders messages in the order they arrive off the socket, and flow 2โ€™s reconnect has no cursor to resume from because there is no โ€œpositionโ€ in a conversation to hold. Deep Dive 1 just made this materially worse: with one server, arrival order at least matched send order for a single sender. Now Aliceโ€™s two messages can be handled by two different servers.

Why this is hard: In a distributed system, thereโ€™s no global clock. Server Aโ€™s timestamp might be 50ms ahead of Server B. Network latency varies. Messages can be retried out of order.

Bad โ€” sort by arrival, or equivalently by the server timestamp the high-level design already stamps on.

This is what we have today. Flow 1 writes the message with whatever timestamp the receiving server put on it, and the client renders in the order frames show up.

The failure is not theoretical and it does not need heavy load to appear. NTP keeps servers within roughly 10-50ms of each other, and Aliceโ€™s two messages are ~1 second apart in human terms but frequently under 50ms apart on the wire when a client flushes a queued send. So:

Concretely: at 1.2M messages/sec, even a 1-in-a-million ordering inversion is more than one visibly scrambled conversation per second, every second.

More detail on the clock problem:

Good - per-conversation monotonic sequence number.

Assign a strictly increasing seqNo per conversation. Every message in a conversation gets the next number in sequence.

Implementation: Redis INCR on key conv_seq:{conversationId}.

Alice sends "Hello"     โ†’ server does INCR conv_seq:alice_bob โ†’ gets 42
Alice sends "How are you?" โ†’ server does INCR conv_seq:alice_bob โ†’ gets 43

Bobโ€™s client sorts by seqNo regardless of arrival order. Even if msg 43 arrives before 42 (network jitter), the UI holds 43 and renders after 42 arrives.

Why Redis INCR? Atomic, single-threaded, sub-ms. Even at 100K messages/sec across all conversations, one Redis cluster handles it because each conversation is an independent key (no contention across conversations).

What about gaps? If Bob receives seqNo 42 then 44 (missed 43), client knows thereโ€™s a gap and requests: โ€œgive me message 43 for this conversation.โ€ Server fetches from the message store.

Great - sequence numbers + client vector clock + multi-device sync.

For apps with multiple devices (phone + web + desktop), ordering gets harder. User sends from phone (seqNo 42), then from desktop (seqNo 43). Both devices need to converge.

The approach (used by WhatsApp, Slack, Facebook Messenger):

  1. Server is the source of truth for sequence numbers. Server assigns seqNo on receipt - NOT the client.
  2. Each device maintains a cursor: lastSyncedSeqNo. On reconnect, device says โ€œgive me everything after seqNo 38โ€ and server sends the delta.
  3. Client embeds lastSeenSeqNo in outgoing messages so the server can detect if the client missed something and proactively push missing messages.
  4. Conflict resolution for near-simultaneous sends from multiple devices: Both get seqNos from the same atomic counter, so theyโ€™re naturally ordered by who hit the server first. No conflict possible at the ordering level.

Tech used in production:


3) The network drops an ack. How do we retry without showing the message twice?

Problem: Flow 1 step 6 relays Bobโ€™s delivered ack back to Alice, and treats the absence of one as nothing at all. Our NFR is zero message loss, which means every un-acked delivery has to be retried โ€” and every retry can duplicate.

In simple terms: The internet is flaky. A message might arrive twice if the โ€˜got itโ€™ confirmation gets lost. We need Bob to see each message exactly once even when the system retries.

Bad โ€” fire and forget, which is what flow 1 does. The Chat Server writes the frame down Bobโ€™s socket and moves on. A TCP write returning success means the bytes reached the kernel buffer, not that Bobโ€™s app processed them, so a client that crashes between receiving and persisting loses the message with no trace. Nobody retries because nobody is tracking that a delivery is outstanding.

At 1.2M messages/sec, a delivery loss rate of even 0.01% is 120 messages lost every second on a requirement that says zero.

Good โ€” retry until acked. Hold each delivered-but-un-acked message and re-push on a timer. That closes the loss hole and immediately opens a duplicate hole: the common failure is not the message being lost, it is the ack being lost, in which case Bob already has the message and gets it again on every retry. Retrying alone converts a loss problem into a duplicate problem.

Great โ€” at-least-once on the wire, deduplicated at the edges. Accept that the network makes exactly-once delivery impossible and make duplicates harmless instead.
๐Ÿ’ก Idempotency = doing the same operation twice has the same effect as doing it once. Here we get it by making every message carry a stable identity that both ends can check against what they have already seen. Learn more โ†’

Flow:

Sender โ†’ Server: message (clientMsgId: "abc")
Server โ†’ Sender: ack (messageId: "msg_1", clientMsgId: "abc")
Server โ†’ Receiver: message (messageId: "msg_1")
Receiver โ†’ Server: delivered ack (messageId: "msg_1")

What if the receiverโ€™s ack is lost? Server retries delivery. Bobโ€™s client sees msg_1 a second time, looks it up in its local DB by messageId, and drops it. The retry is invisible to him.

What if the senderโ€™s send is retried? Aliceโ€™s client reuses the same clientMsgId, so the server checks it against a short-lived dedupe cache and returns the original messageId without storing a second copy. This is why clientMsgId is generated on the client: only Aliceโ€™s device knows that her retry is the same logical send.

Result: at-least-once on the wire, exactly-once from the userโ€™s point of view. This is also what makes Deep Dive 5โ€™s crash-and-redeliver safe โ€” a redelivered fan-out batch re-pushes messages some members already have, and their clients discard them.


4) How does a message sent from a phone reach the laptop that was asleep?

Problem: Flow 2 step 6 waves its hands: the reconnecting client asks for โ€œeverything it does not already have.โ€ Deep Dive 2 gave us the sequence number that makes that expressible. This is where we cash it in โ€” and generalise it, because a user has several devices and each one is at a different position.

Note this is formally below the line in our NFRs, but it falls out of the sequence number nearly for free, and flow 2 does not actually work without it.

In simple terms: You send a message from your phone. When you open WhatsApp on your laptop 5 minutes later, that same message should be there. Every device needs to know its own place in the conversation.

Bad โ€” โ€œsend me the last N messagesโ€ on every reconnect, which is the only thing flow 2 could do before Deep Dive 2. The client pulls a fixed window and diffs it locally against what it has. This is wrong in both directions: if Bob was away for a week, N is too small and he silently loses the middle of the conversation; if he was away for 30 seconds, we just re-sent N messages to answer a question whose real answer was โ€œtwo.โ€ At ~500K reconnect fetches/sec that is a large amount of bandwidth spent re-sending data the client already had, and it still has a correctness hole.

Good โ€” one cursor per user. Track lastSyncedSeqNo per conversation for the account and send the delta after it. Correct for one device, wrong for three: whichever device syncs first advances the shared cursor, and the other two never receive those messages at all.

Great โ€” a cursor per device, over a server-authored ordered log. The server owns the sequence; each device is an independent reader with its own position.

This is the โ€œordered logโ€ model (Facebook Iris). The server is the source of truth; clients are materialized views with a cursor. It is also what makes an offline queue worth adding: not to hold the messages, which the Message Store already does durably, but as a compact per-user index of which conversations moved while a device was away, so reconnect does not have to ask that question of every conversation the user is in.


5) A 500-person group gets one message. Do we write it 500 times or read it 500 times?

Problem: Flow 3 delivers a group message by looping over the member list inside the senderโ€™s request. That is correct and it makes the sender pay for the size of the group.

In simple terms: When you send a message to a 500-person group, should we write 500 copies - one per member - or write one copy and let each member fetch it? And either way, the person who hit send should not be waiting while it happens.

Bad โ€” the in-request loop from flow 3. Two failures, and the latency one is arithmetic.

Latency. A 256-member group means up to 256 socket writes plus up to 256 calls to FCM or APNs, each an external HTTPS round trip at roughly 50ms. Done sequentially that is 12.8 seconds before Aliceโ€™s client is acked, against a 500ms P99. Firing them in parallel helps the mean and not the tail: the sender still waits for the slowest of 256 external calls, and now we hold 256 concurrent outbound sockets per group send.

Durability of delivery. If the process dies at member 200, members 201-256 are never delivered and nothing records that they were skipped. The message is safely stored โ€” flow 1 stores before delivering โ€” so it is not lost. It is silently undelivered, which for those 56 people is indistinguishable from lost. There is no retry, because there is no record of what still needs retrying.

Good โ€” hand the fan-out to a background worker. Ack the sender as soon as the message is stored, then let a worker do the member loop. Send latency becomes constant regardless of group size, which is the whole point. But an in-process background worker still loses the remaining members if it dies, so the durability hole is untouched โ€” we have only moved it off the senderโ€™s critical path.

Great โ€” a durable log between accepting and delivering. Publish one message.fanout event to an event bus (Kafka / Redpanda / Pub-Sub) and let a pool of Fan-out Workers consume it.
๐Ÿ’ก The event bus is not here to be fast. It is here so that โ€œwho still needs this messageโ€ survives a worker crash. Learn more โ†’

Mechanism:

  1. Chat Service stores the message, publishes one event keyed by conversationId, and acks the sender. Sender latency is now independent of group size
  2. Keying by conversationId also buys ordering for free: one partition per conversation means one consumer handles that groupโ€™s messages in sequence, so the Deep Dive 2 sequence numbers stay monotonic per group
  3. A Fan-out Worker consumes the event, reads the member list, and runs flow 1 or flow 2 per member โ€” socket push if connected, Push Service if not
  4. The worker commits its offset only after the batch completes. Crash mid-batch and the event is redelivered, so members get re-attempted rather than skipped. This is at-least-once, which is exactly what Deep Dive 3โ€™s client dedupe already handles
  5. A rate limiter caps how much of the worker pool one conversation may occupy, so a 100K-member channel cannot starve everyone elseโ€™s groups

And now the storage question the heading asks, which is separate from the delivery question and has a different answer at each size:

ย  Push model (write 500 times) Pull model (read 500 times)
Writes per message One copy per member. A 256-member group at 1,000 messages/day is 256K writes/day for one group One copy, regardless of member count
Reads Fast โ€” each user reads only their own inbox partition Each read merges across the conversations the user belongs to
Breaks when Membership gets large. A 100K-member channel is 100K writes for one message Membership gets small and numerous โ€” read merge cost grows with conversation count

Hybrid (what WhatsApp and Discord do): push for small groups where fan-out is bounded and reads should be trivial, pull for large channels where the write cost is the binding constraint. The threshold sits around a few hundred members and is worth tuning rather than guessing. Note that our high-level design already chose pull, storing one copy per conversation โ€” so the work here is adding the push path for small groups, not replacing what we have.

Final Architecture

flowchart TD
    CLIENTS["Mobile and Web Clients"]:::client
    LB["Load Balancer<br/>sticky by userId"]:::edge
    WS["WebSocket Servers<br/>Netty edge tier"]:::service
    CHAT["Chat Service"]:::service
    REG[("Connection Registry<br/>Redis")]:::data
    STORE[("Message Store<br/>Cassandra")]:::data
    OFFLINE[("Offline Queue<br/>Redis sorted set")]:::data
    K["Kafka<br/>fan-out and events"]:::async
    FAN["Fan-out Workers"]:::service
    PUSH["Push Service"]:::service
    MEDIA[("S3 and CDN<br/>media")]:::data
    FCM["FCM and APNs"]:::external

    CLIENTS -->|"Open WebSocket"| LB
    CLIENTS -->|"Presigned upload"| MEDIA
    LB -->|"Sticky route by user"| WS
    WS -->|"Forward to chat logic"| CHAT
    CHAT -->|"Lookup receiver server"| REG
    CHAT -->|"Persist message"| STORE
    CHAT -->|"Queue for offline user"| OFFLINE
    CHAT -->|"Publish group fan-out"| K
    K -->|"Process group delivery"| FAN
    FAN -->|"Push to online members"| WS
    CHAT -->|"Trigger push alert"| PUSH
    PUSH -->|"Deliver via FCM APNs"| FCM

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef edge fill:#1e3a5f,stroke:#60a5fa,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef async fill:#AB47BC,stroke:#4A148C,color:#fff
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
    classDef external fill:#4a1942,stroke:#f472b6,color:#e2e8f0

How it works end-to-end:

  1. Client opens WebSocket โ€” connects through Load Balancer (sticky by userId) to a WebSocket Server
  2. Sender sends message โ€” WebSocket Server forwards to Chat Service
  3. Chat Service persists message โ€” writes to Cassandra (Message Store) with a per-conversation sequence number
  4. Connection Registry checked โ€” Redis lookup finds which WebSocket Server holds the recipient
  5. Kafka fan-out for groups โ€” message event published to Kafka, Fan-out Workers push to each memberโ€™s WebSocket Server
  6. Recipient online โ€” message delivered in real-time through their WebSocket connection
  7. Recipient offline โ€” message queued in Redis sorted set (Offline Queue) and push notification sent via FCM/APNs
  8. Recipient reconnects โ€” drains Offline Queue in order, syncs from last seen sequence number

Summary

Decision Choice Why
Transport WebSocket Real-time bidirectional, sub-second delivery
Message store Cassandra Partition per conversation, append-only, handles billions
Connection registry Redis Sub-ms lookup of โ€œwhich server has user Xโ€
Offline delivery Redis sorted set + push notification Ordered drain on reconnect
Group fan-out Kafka โ†’ workers Async, retryable, doesnโ€™t block sender
Ordering Per-conversation sequence number Simple, no clock dependency
Delivery guarantee At-least-once + client dedupe Zero message loss, no duplicates visible to user
Multi-device Pull sync with seqNo cursor Ordered log model (Facebook Iris)

Key Technologies

Term What it is
WebSocket A persistent two-way connection between client and server. Unlike HTTP (request โ†’ response โ†’ done), WebSocket stays open so the server can push messages to the client anytime.
Cassandra A distributed NoSQL database optimized for fast writes. Stores data across many machines. Perfect for append-only message logs.
Kafka A distributed event streaming platform. Producers write events, consumers read them. Used here to decouple message sending from delivery fan-out.
Redis In-memory key-value store (< 1ms reads). Used here for connection registry (which user is on which server) and offline message queues.
FCM / APNs Firebase Cloud Messaging (Android) and Apple Push Notification service (iOS). How you send push notifications to phones when the app is closed.
Sequence Number A monotonically increasing integer per conversation. Guarantees message ordering regardless of clock differences between servers.
Store-and-Forward Pattern where the server stores a message durably first, then delivers it when the recipient is available. Ensures zero message loss.
Fan-out Delivering one message to multiple recipients (group chat). โ€œFan-out on writeโ€ = copy to each inbox. โ€œFan-out on readโ€ = store once, each client fetches.

Whatโ€™s Expected at Each Level

This section helps you calibrate your depth. You donโ€™t need to cover everything - just know whatโ€™s expected for your level.

Mid-level

Design basic 1:1 messaging with a server relaying messages. Propose WebSocket for real-time delivery. Understand offline message storage and why polling is wasteful. With prompting, discuss how to handle group messages by fanning out to multiple recipients.

Senior

Propose Cassandra for message storage (partition by conversation). Explain connection-level routing - how does a message find the right WebSocket server? Discuss read receipts, message ordering guarantees (per-conversation sequence numbers), and offline delivery queues. Articulate why eventual consistency is acceptable for message delivery.

Staff+

Address end-to-end encryption key exchange (Signal protocol double-ratchet), multi-device sync with ordered-log cursors, and message fan-out for large groups (1000+ members) using the hybrid push/pull model. Discuss graceful degradation when the chat service is overloaded (backpressure on WebSocket connections). Cover message retention policies and GDPR right-to-deletion across replicated stores.


๐ŸŽฏ Key Takeaways



Understand the building blocks used in this design:

Discussion

Newest first
You

Free system design + DSA prep. If it helped you crack an interview, consider supporting.

SensAI SensAI
Beta
Listening...
Tap mic to stop voice mode

Shape what we build next

Every piece of feedback is read by the team and directly influences our roadmap.

What type of feedback?

Install SystemCraft

Add to your home screen for instant access, offline reading, and a distraction-free experience.

Offline reading Faster loads No browser tabs App-like feel

Unlock AI Features

One click to activate - no payment, no credit card. Just sign in and you're in.

AI code review and hints
SensAI chat assistant
AI mock interviews
Whiteboard analysis
100% free during early access