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:
- Polling is wasteful - 500M users polling every 5s = 100M QPS of mostly-empty responses. Massive cost, terrible latency.
- 5-second delay feels laggy - real-time chat needs sub-second delivery.
- Single DB for all messages - billions of messages/day, one table collapses.
- No offline handling - if receiver is offline when message arrives, when do they get it?
- No ordering guarantee - if two messages arrive out of order at the DB, display is wrong.
- Group messages multiply the problem - 256-member group = 256 deliveries per message.
Prior Art Weโre Drawing From
- WhatsApp Architecture (InfoQ) - Erlang-based, 2M connections per server, XMPP-derived protocol, store-and-forward for offline delivery.
- Facebook Messenger Iris - ordered log storage (like Kafka) per conversation. Messages appended to a per-user ordered log. Clients sync via sequence numbers.
- Discord How Messages Are Stored - migrated from MongoDB to Cassandra to ScyllaDB. Partition per channel + bucket.
- Signal Protocol - end-to-end encryption with double-ratchet. Pre-keys for offline delivery. The gold standard for E2E chat encryption.
- Slack Real-Time Messaging - WebSocket connections for real-time, application-level edge cache (Flannel) for fast channel hydration.
Functional Requirements
Core:
- Users can send messages (text) to another user in real-time (1:1 chat).
- Users can create groups and send messages to all group members.
- Messages are delivered reliably even if the recipient is offline (store-and-forward).
Below the line:
- Read receipts, typing indicators
- Media messages (images, video, voice)
- End-to-end encryption
- Message search, reactions, threads
- Voice/video calling
Non-Functional Requirements
Core:
- Real-time delivery - P99 < 500ms for online-to-online message delivery.
- Reliability - zero message loss. Once the server acks, the message WILL be delivered eventually.
- Ordering - messages within a conversation appear in send order.
- Scale - 500M DAU, 100B messages/day and 50M concurrent connections (WhatsApp scale).
Below the line:
- Sub-100ms delivery latency
- Exactly-once delivery (at-least-once + client-side dedupe is acceptable)
- Multi-device sync (web + mobile + desktop)
Scale Estimation (Back-of-Envelope)
- Users: 500M DAU, 50M concurrent connections at peak
- Write QPS: 100B messages/day = ~1.2M messages/sec sustained, ~3M/sec at peak
- Read QPS: ~500K message fetches/sec (history sync + offline drain on reconnect)
- Storage: 100B messages/day ร ~200 bytes = ~20TB/day, so ~7.3PB/year raw and ~2PB/year after compression
- Bandwidth: ~50 Gbps aggregate at peak (payload plus WebSocket framing, heartbeats and acks)
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
- User - identified by phone number or userId. Has online/offline status.
- Conversation - a 1:1 or group thread. Has a unique
conversationIdand list of participants. - Message - text content with
messageId,senderId,conversationId,timestamp,status(sent/delivered/read). - Connection - a live WebSocket session mapping
userId โ serverId:connectionId.
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:
- 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 โ connectiontable 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 โ - Message Store (Cassandra) - permanent storage for every message, partitioned by
conversationIdso 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:
- Sender types โHey, are you free tonight?โ and hits send โ the frame travels over their already-open WebSocket to the Chat Server
- 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 - Chat Server looks up the receiver in its in-memory connection map โ found, Bob is connected
- It writes the message frame down Bobโs socket. It appears on his screen
- Bobโs app replies with a
deliveredack - 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:
- One server. 50M concurrent connections do not fit on one machine, and the in-memory connection map dies with the process, so a deploy disconnects everybody. The moment there are two servers, step 3 stops working: Aliceโs server has no idea Bobโs socket lives on another host. That is Deep Dive 1.
- No ordering. Nothing in steps 1-6 assigns an order. Bob renders messages as they arrive off the socket, which is not the same as the order Alice sent them. That is Deep Dive 2.
- Acks are optimistic. Step 6 tells Alice the message was delivered, but if step 4โs frame was lost in flight we would never know, and if Bobโs ack in step 5 is lost we will retry and show him the message twice. That is Deep Dive 3.
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):
- 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. - 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:
- Server sends a silent data message via FCM - just a โwake up, you have a new messageโ signal with sender ID and message reference
- FCM wakes up the appโs background process on the device
- The app connects to the server, pulls the encrypted message, and decrypts it locally on the device
- 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:
- Alice sends the message exactly as in flow 1
- Chat Server persists it to the Message Store first, as always โ store, then deliver
- It looks up Bob in the connection map and does not find him. Bob is offline
- Push Service asks FCM or APNs to wake Bobโs device
- Bobโs phone buzzes with โNew message from Aliceโ
- Whenever Bob next opens the app, his client reconnects and asks for everything it does not already have in each conversation
- 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.
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:
- Sender sends a message to group
conv_123, which has 256 members - Chat Server stores ONE copy with partition key
conv_123, not 256 copies - It reads the 256 member IDs from
conversation_members - It loops over them. For each member found in the connection map, write the frame down their socket
- For each member not found, call Push Service โ flow 2, once per absent member
- 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:
- The sender waits for everyone. Send latency is now proportional to group size. A 256-member group means 256 socket writes and up to 256 push calls, each an external HTTPS round trip to FCM, before Aliceโs client hears back. That blows the 500ms P99 for the sender on a requirement that has nothing to do with the sender.
- A crash halfway through loses the rest. If the process dies at member 200, members 201-256 never get delivered and nothing anywhere records that fact. The message is safely stored, so it is not lost โ but it is silently undelivered, which for the 56 people who never see it is the same thing.
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:
- Sender sends message โ Chat Service calls
INCR seq:{conversationId}โ gets seqNo 42 - Message stored in Cassandra with
(conversation_id, seqNo=42)as the clustering key - Fan-out workers deliver messageId + seqNo to each online member via WebSocket
- Receiverโs client inserts message at correct position by seqNo (not arrival time)
- 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 - All devices sync via the same seqNo cursor โ each device tracks
lastSyncedSeqNoand 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:
- Each idle WebSocket = ~10KB RAM (file descriptor + small read/write buffers)
- 2M connections ร 10KB = 20GB RAM (fits in a 64GB server)
- Netty can handle 1-2M connections per JVM instance on modern hardware
- Goโs goroutines achieve similar density (goroutine = ~4KB stack vs Java thread = 512KB)
Tech used in production:
- WhatsApp: Erlang/OTP (lightweight processes, similar to goroutines - famously ran 2M connections per server)
- Discord: Elixir/Erlang on the gateway, Rust for hot paths
- Slack: Java + Netty for WebSocket gateway (project โFlannelโ)
- Signal: Java + Netty
- WeChat: C++ custom framework
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โ:
- ONLY manages TCP/WebSocket connections, TLS handshake, and heartbeat pings
- Does NO business logic - just holds open connections and passes messages through
- Extremely lightweight: each connection costs ~10KB (just a file descriptor + buffer)
- Can hold 2M+ connections per node because itโs barely doing anything per connection
- Built with: Envoy proxy, HAProxy, custom Go/Rust services, or Netty with minimal handlers
Logic Tier (message processing) - the โbrainโ:
- Receives raw messages from edge tier via internal gRPC
- Handles all business logic: routing, persistence to Cassandra, fan-out to group members, push notifications
- Stateless - scales horizontally based on message throughput
- Doesnโt hold any WebSocket connections - just processes and responds
How a message flows through both tiers:
- Aliceโs phone is connected to Edge Server #3 via WebSocket
- Alice sends โHey Bobโ โ Edge Server #3 receives the raw bytes
- Edge Server #3 forwards to Logic Tier via internal gRPC: โmessage from userId=alice, payload=Hey Bobโ
- Logic Tier stores in Cassandra, then checks Connection Registry: โBob is on Edge Server #7โ
- Logic Tier sends to Edge Server #7: โdeliver this to Bobโs WebSocket connectionโ
- Edge Server #7 pushes the message down Bobโs WebSocket
- 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:
- Adding more connections = adding cheap, lightweight edge nodes (no processing overhead)
- A slow DB write in Logic Tier doesnโt block Edge Tier from handling new connections/pings
- If an edge node crashes: clients reconnect to another edge node. No messages are lost (Logic Tier handles durability separately)
- During idle hours (3 AM): connections exist but messages are rare. Edge handles the load efficiently, Logic Tier is mostly idle
Real-world implementations:
- WhatsApp: Erlang nodes at edge, backend services for routing/storage
- Discord: โGatewayโ servers (Elixir) hold connections, โGuildโ servers handle message logic
- Slack: โFlannelโ is their edge/cache layer, backend services do the real work
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:
- Aliceโs โHelloโ lands on Server A whose clock reads T=1000. โHow are you?โ lands on Server B whose clock reads T=999. Bob sees the reply before the question, permanently, because that is what we persisted.
- Even on a single server, two messages in the same millisecond have no defined order at all.
- And flow 2 has no way to express โeverything after X,โ so a reconnecting client cannot ask for a delta.
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:
- Clock skew between servers (NTP syncs every few seconds, drift is 10-50ms)
- Aliceโs โHelloโ hits Server A at T=1000, โHow are you?โ hits Server B whose clock reads T=999. Bob sees them reversed.
- Even on one server, if two messages arrive in the same millisecond, order is random.
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):
- Server is the source of truth for sequence numbers. Server assigns seqNo on receipt - NOT the client.
- Each device maintains a cursor:
lastSyncedSeqNo. On reconnect, device says โgive me everything after seqNo 38โ and server sends the delta. - Client embeds
lastSeenSeqNoin outgoing messages so the server can detect if the client missed something and proactively push missing messages. - 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:
- WhatsApp: Server-assigned message IDs + per-chat ordering. Each message has a globally unique ID + per-conversation sequence.
- Slack: Uses a
ts(timestamp) as the unique message ID within a channel. Server-generated, monotonically increasing per channel. Format:1234567890.123456. - Discord: Snowflake IDs (time-based, globally unique). Messages sorted by Snowflake ID which is inherently time-ordered since timestamp is the most significant bits.
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.
- Each conversation has a
maxSeqNo - Each device tracks its own
lastSyncedSeqNoper conversation, so devices cannot advance each otherโs position - On app open, the device sends
GET /conversations/:id/messages?after=lastSyncedSeqNoand gets exactly the delta - Live messages arrive over the WebSocket and advance the deviceโs cursor as they are persisted locally
- A device that has been off for a month and one that has been off for a minute use the identical code path โ the only difference is the size of the delta
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:
- Chat Service stores the message, publishes one event keyed by
conversationId, and acks the sender. Sender latency is now independent of group size - Keying by
conversationIdalso 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 - 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
- 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
- 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:
- Client opens WebSocket โ connects through Load Balancer (sticky by userId) to a WebSocket Server
- Sender sends message โ WebSocket Server forwards to Chat Service
- Chat Service persists message โ writes to Cassandra (Message Store) with a per-conversation sequence number
- Connection Registry checked โ Redis lookup finds which WebSocket Server holds the recipient
- Kafka fan-out for groups โ message event published to Kafka, Fan-out Workers push to each memberโs WebSocket Server
- Recipient online โ message delivered in real-time through their WebSocket connection
- Recipient offline โ message queued in Redis sorted set (Offline Queue) and push notification sent via FCM/APNs
- 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
- WebSocket for real-time delivery - persistent connection, server pushes
- Cassandra for message storage - partitioned by conversation for fast reads
- Store-and-forward for offline users - deliver when they reconnect
- Connection registry in Redis routes messages to the right WebSocket server
Related Designs
- Notification System - similar multi-channel delivery + WebSocket patterns
- Twitter Feed - fan-out and real-time updates
- Stock Broker - Kafka event streaming + exactly-once semantics
Related Concepts
Understand the building blocks used in this design:
- WebSockets vs SSE โ โ persistent connections deliver messages in real time both ways
- Message Queues โ โ buffers and routes messages between senders and recipients
- Fan-Out Patterns โ โ delivers a single group message to every member of the conversation
- Consistent Hashing โ โ maps each user to a connection server so messages find the right socket
Discussion
Newest first