Course Content
System Design Interview
31 sections · 71 lessons
Chat System: scope, scale, connections and the stateful tier
The prompt: "Design a chat application."
Attempt it for 45 minutes on paper before reading. Draw the connection path, write down what a message row looks like, and try to answer "which server does the recipient's message go to?" That last question is the one this section is really about.
This lesson covers the part of the design that is unlike anything before it: scoping, the estimate that points at connections rather than messages, how the server pushes to the client, and what it means for a fleet of 200 servers to each hold live, stateful connections.
Why chat is different from everything before it
Every design so far has been request-and-response: a client asks, a server answers, and the connection closes. Chat inverts that. The interesting event — someone sends you a message — originates on the server, and the client has to learn about it without asking.
That single inversion produces the two hard parts of the design, both covered later in this lesson: how the server pushes, and the fact that a server holding a live connection is stateful in a world built on stateless web tiers.
The five questions
1. One-to-one, group, or both? One-to-one chat is a two-person channel and is the simple case. Group chat is fan-out, and the group size limit changes the design: 100 members is a loop, 100,000 members is a broadcast system. Ask for the cap.
2. How long is history retained? Forever changes the storage design completely — it is the difference between a rolling 30-day buffer and a permanently growing archive. Assume forever unless told otherwise, and size it.
3. Which features beyond text? Read receipts, typing indicators, presence, and media each add traffic. Typing indicators in particular can generate more events than messages do.
4. Is end-to-end encryption in scope? End-to-end encryption means the server stores ciphertext it cannot read, which removes server-side search, server-side moderation, and easy multi-device sync. If it is in scope, say what it costs rather than treating it as a checkbox; Follow-ups covers the four consequences.
5. What is the concurrent-connection count? Not daily active users — the number of people connected at the same moment. That number, and not message volume, sizes the connection tier.
Scope
In: one-to-one and group messaging, delivery and ordering, presence, offline delivery. Out: voice and video calls (a different system built on real-time media transport), the recommendation of who to chat with, and the client user interface.
Requirements and scale
Functional: send a message to a person or a group; receive messages in near real time; read history; see delivery and read state; see who is online; receive messages sent while offline.
Non-functional: message delivered to an online recipient in under 500 ms end to end; messages within a conversation appear in a consistent order for everyone; no message is ever lost once acknowledged; the system stays available during a partial failure, with delayed delivery preferred over refused sends.
The estimate
Invented figures for a large messaging product. A day is 100,000 seconds.
| Quantity | Assumption | Result |
|---|---|---|
| Daily active users | 50 million | — |
| Messages sent per user per day | 40 | 2B messages/day |
| Average message rate | 2B ÷ 100,000 | 20,000 messages/s |
| Peak message rate | 3× | 60,000 messages/s |
| Average message size | 300 bytes stored | — |
Storage: 2B × 300 bytes = 600 GB/day, or about 220 TB/year, growing forever if history is retained. That rules out a single relational cluster and points at a horizontally partitioned store (Message storage and ordering).
The connection number
Assume 30% of daily active users are connected at peak:
50M × 0.30 = 15 million concurrent connections
Now size the tier. A connection costs socket buffers, transport-layer state, and per-user application state — call it roughly 10 KB of kernel and application memory each, though the real figure depends heavily on buffer tuning and should be measured rather than assumed. A well-tuned server holding 100,000 connections therefore needs about 1 GB for connection state alone, plus headroom.
15,000,000 ÷ 100,000 per server = 150 connection servers, before redundancy.
Round to 200 for failure headroom and deploy capacity. That is a real, ordinary fleet — and every one of those 200 servers is stateful, which is the subject of the last part of this lesson.
Bandwidth
At 60,000 peak messages per second and, say, 500 bytes on the wire including framing, the inbound message traffic is 30 MB/s. Outbound is larger because group messages fan out, but it is still tens to low hundreds of megabytes per second. Bandwidth is not the constraint here — unlike Section 16 (Design YouTube), where it is the entire design.
How the server pushes to the client
The client needs to learn about an event it did not request. There are four ways, they are not equivalent, and the comparison is a standard interview probe. Communication patterns, in Section 5, covers these patterns in general; here they are decided for chat specifically.
The four options
Polling. The client asks "anything new?" on a timer. Simple, works everywhere, and wasteful. At a 2-second interval for 15 million connected clients:
15,000,000 ÷ 2 = 7.5 million requests per second
against a real message rate of 60,000 per second. More than 99% of those requests return nothing. And the average message still waits 1 second before the client asks for it.
Long polling. The client asks and the server holds the request open until there is something to send, or a timeout of 30 to 60 seconds fires. Message latency drops to near zero. The request rate falls to roughly one request per client per timeout period: 15M ÷ 45 s ≈ 330,000 requests per second, a 20× improvement on polling. The costs are real: the server holds an open request per client anyway (so it is already stateful), each message requires a full new request afterwards, and it is one-directional.
Server-sent events. One long-lived HTTP response that the server writes into repeatedly. Efficient, automatically reconnects, and travels over ordinary HTTP. But it is server-to-client only, so sending a message needs a separate HTTP request. For a system that is symmetric by nature, that asymmetry is awkward.
WebSockets. One connection, upgraded from HTTP, carrying frames in both directions for as long as it lives. Send and receive share it. Per-message overhead is a few bytes of framing instead of a full set of HTTP headers.
The comparison, and the recommendation
| Latency | Requests/s at 15M clients | Bidirectional | Server state | |
|---|---|---|---|---|
| Polling (2 s) | up to 2,000 ms | 7,500,000 | n/a | none |
| Long polling (45 s) | ~0 ms | 330,000 | no | one open request per client |
| Server-sent events | ~0 ms | 15M open responses | no | one open response per client |
| WebSockets | ~0 ms | 15M open sockets | yes | one open socket per client |
Recommend WebSockets, with the reason stated: chat is symmetric and high-frequency, so a single bidirectional connection removes both the request overhead of long polling and the send-path asymmetry of server-sent events.
Then say the honest part. The state cost is identical for long polling, server-sent events, and WebSockets — all three pin a client to a specific server. Choosing WebSockets does not create the stateful problem below; it makes it explicit.
The stateful problem
Every architecture in this course so far has relied on stateless web servers: any request can go to any machine, and a dead machine costs nothing but a retry (Load balancers and stateless web tiers). A chat server holds live connections, so that assumption is gone. This is the deep dive, and the one new idea this section owns.
The problem in one sentence
Aditi is connected to chat server 47. Ravi, connected to server 112, sends Aditi a message. Server 112 has the message and no way to reach Aditi, because Aditi's socket lives inside a different process on a different machine. (Aditi and Ravi are invented examples used throughout this section.)
The session registry
The fix is a lookup table from user to server, written when a connection opens and deleted when it closes.
Notice step 3: without the registry, server 112 has the message and no way to reach a socket that lives in another process on another machine.
Two implementation choices worth naming:
- A key-value store with a short time-to-live.
session:u_aditi → server-47, refreshed by a heartbeat every 10 seconds with a 30-second expiry. Simple, fast, and self-healing: a crashed server's entries expire on their own. - A coordination service such as ZooKeeper or etcd, with ephemeral entries tied to the server's own session. Stronger consistency, and it disappears the instant a server dies.
Recommend the key-value store with a time-to-live for this use case. Registry reads happen on every message — 60,000 per second at peak — and a coordination service is built for consistency at low write rates, not for a read-heavy hot path. The stale-entry window a time-to-live creates is tolerable, because the sending server discovers the truth when it tries to forward and the target refuses.
Routing between servers
Server 112 does not open an HTTP connection to server 47 for every message. Two better shapes:
- Direct forwarding over a persistent internal connection mesh. Lowest latency, but a 200-server fleet means each server keeps up to 199 internal connections.
- A per-server topic in a distributed log, where server 47 consumes the topic for messages destined to its connected users. Higher latency (a few milliseconds), far simpler operationally, and it buffers when a target server is briefly busy.
Recommend the log for a first design and mention direct forwarding as the latency optimisation, because the log also gives you the offline path for free — Follow-ups uses it.
When a server dies
Server 47 crashes with 100,000 connections on it.
- Those 100,000 clients see their socket close.
- They reconnect — and this is the dangerous moment. One hundred thousand simultaneous reconnects with authentication is a self-inflicted denial of service. Clients must reconnect with exponential backoff and jitter (Reliability patterns), spreading the reconnect storm over tens of seconds.
- Registry entries for those users expire within 30 seconds. Until then, messages routed to server 47 are undeliverable — so the sending server must fall back to the offline path (store the message in the recipient's inbox) rather than dropping it.
- Messages sent during the gap are picked up when the client reconnects and syncs from its last received message identifier.