Problem & requirements
Chat apps vary widely, from Slack-style team tools to Discord-style gaming chat to WhatsApp-style personal messaging, so pin down scope first. Assume one-on-one chat with low latency, small group chats of up to 100 members, online presence indicators, text-only messages up to 100,000 characters, multi-device login for the same account, push notifications for offline users, and 50 million daily active users. End-to-end encryption can be mentioned as a follow-up.
The defining property is that the server must push messages to recipients the moment they arrive. Ordinary request-response HTTP is built for the client to initiate, so the first design question is how the server reaches the client.
Back-of-the-envelope estimation
With 50M DAU sending 40 messages each per day, the system handles 2 billion messages a day, or 2,000,000,000 / 86,400 ≈ 23,000 messages per second, with peaks perhaps three times that. At about 100 bytes per average message, that is 200 GB of new history per day, or 73 TB per year, which grows forever because chat history is usually kept.
Connections are the other dimension. If one server holds 1M concurrent WebSocket connections at roughly 10 KB of memory each, that is 10 GB per server, so a modern machine could in theory host every connection for a much smaller app. Putting everything on one server is still a bad idea because it becomes a single point of failure.
Polling, long polling and WebSocket
Polling has the client ask 'anything new?' at fixed intervals. It is simple but wasteful, because most answers are 'no' and each request costs a round trip. Long polling holds the request open until a message arrives or a timeout fires, which reduces empty responses, but the sender and receiver may hit different servers, the server cannot easily tell when a client disconnected, and inactive users still cycle connections.
WebSocket starts as an HTTP request and is upgraded to a persistent, bidirectional connection over the same ports 80 and 443, so it usually passes through firewalls. Once open, the server can push at any time. Using WebSocket for both sending and receiving simplifies the design, at the cost of managing many long-lived connections on the server side.
High-level design: stateless, stateful and third-party parts
Most features, such as sign-up, login, profiles, and group management, are ordinary stateless HTTP services behind a load balancer. One important stateless piece is service discovery (for example, Apache ZooKeeper), which picks the best chat server for a client based on location and load and hands back its address.
The chat service is stateful: each client keeps a persistent connection to one chat server and does not switch unless that server becomes unavailable. A presence server tracks online status. Push notifications reach users with no open connection, and an API server layer handles everything non-real-time.
Storage and message IDs
Generic data such as user profiles and friend lists fits a replicated, sharded relational database. Chat history is different: the volume is huge, recent messages are read far more than old ones, users still need random access for search and jump-to-message, and the read-to-write ratio is roughly 1:1 for one-on-one chat. Key-value stores such as HBase or Cassandra suit this because they scale horizontally, offer low-latency access, and handle the long tail of old data better than a relational index. Facebook Messenger used HBase and Discord used Cassandra.
Messages need an order. Timestamps are unreliable because two messages can share one. Each message needs an ID that is unique and sortable by time. A global ID generator like Snowflake works, but a cheaper approach is a local sequence number per channel, since ordering only needs to hold within one conversation, not across all of them.
- 1:1 table:
message_id(primary key),message_from,message_to,content,created_at. - Group table: composite key
(channel_id, message_id), withchannel_idas the partition key.
Message flows and multi-device sync
In a one-on-one flow, User A sends a message to Chat server 1, which obtains a message ID and puts the message on a message sync queue. The message is stored in the key-value store. If User B is online, it is forwarded to the chat server holding B's connection, which pushes it over WebSocket; if B is offline, the push notification servers alert B's devices.
For small groups, the sender's message is copied into each recipient's own inbox queue. This makes reads simple, since each client only checks one inbox, and it is cheap when groups have at most around 100 members. For huge groups the copy cost becomes prohibitive, which is why large channels use a different model.
With multiple devices, each device remembers cur_max_message_id, the newest message ID it has seen. On reconnect it fetches messages where the recipient is the logged-in user and the ID is greater than that value. Because IDs are ordered, each device catches up independently.
cur_max_message_id, so after a reconnect it asks only for messages above that.Online presence
When a user connects, the presence server records the status and a last_active_at timestamp in a key-value store. Explicit logout flips it to offline. Network flakiness makes disconnect events unreliable, especially on mobile, so treating every brief drop as offline would make status flicker constantly.
Instead, clients send a heartbeat every few seconds. If no heartbeat arrives within a window, say 30 seconds, the user is marked offline. To propagate status changes, the presence server publishes to a pub/sub channel per friend pair, and friends subscribe to it. This is cheap for small friend lists or groups, but for very large groups presence is better fetched on demand, for example when a user opens the group.
Trade-offs, failures & wrap-up
If a chat server dies, its clients lose their connections; service discovery hands them a new server and they reconnect, using cur_max_message_id to catch up. Messages that fail to send are retried from the queue. Separating the stateful chat servers from the stateless API tier keeps failures contained.
Extensions worth naming include media messages, which add compression, cloud storage, and thumbnails; end-to-end encryption, where only the endpoints hold keys; client-side caching of recent messages; and geographically distributed caches to improve load time. A strong answer justifies WebSocket, the key-value store, and per-channel IDs rather than just naming them.
Key numbers
Key terms
- Long polling
- A technique where the server holds a client's request open until data is available or a timeout expires.
- WebSocket
- A protocol that upgrades an HTTP connection into a persistent, full-duplex channel.
- Stateful service
- A service whose clients stay bound to a specific server because it holds per-connection state.
- Service discovery
- A component that tells clients which server instance to connect to based on load and location.
- Message sync queue
- A per-user inbox that holds messages waiting to be delivered to that user's devices.
- cur_max_message_id
- The newest message ID a device has seen, used to fetch only newer messages on reconnect.
- Heartbeat
- A periodic signal from a client proving it is still online.
- Pub/sub channel
- A messaging pattern where publishers broadcast events to all subscribers of a topic.
Common mistakes
- Proposing plain polling without discussing its wasted requests and latency.
- Using wall-clock timestamps as the sole ordering key for messages.
- Copying every message into every member's inbox for groups of thousands.
- Marking users offline on every brief network drop instead of using heartbeats.
- Forgetting that a chat server crash breaks connections and how clients recover.
- Storing petabytes of chat history in a single relational table without sharding.
Further study
- How Discord stores billions of messages (Discord engineering blog)
- The Underlying Technology of Messages (Facebook engineering, HBase)
- RFC 6455 The WebSocket Protocol
- Apache ZooKeeper
- Signal Protocol (Double Ratchet)