General

How to Design a Chat System That Scales to Millions of Users

TL;DR Here's the thing most tutorials miss: Server-Sent Events (SSE) is a third option between long polling and WebSockets. SSE is HTTP-based (works through standard HTTP/2 multiplexing), server-push only (client-to-server still uses regular HTTP requests), and has automatic reconnection built in.

Read this article as text (accessible version)
System Design · Real-Time Architecture · 2026 · 4,000 words · 4 interactive labs · April 2026

How to Design a
Chat System
That Scales to
Millions of Users

WhatsApp processes over 100 billion messages a day. Every message that arrives instantly, every green dot that appears correctly, every notification that wakes up a sleeping phone — each one is the product of deliberate architectural decisions made by engineers who understood the tradeoffs. Here's how to design the same system from scratch.

Read the Deep Dive ↓ Open the Lab 💬 WebSocket + HTTP hybrid Inbox pattern for reliability Fan-out for group messages Heartbeat presence Horizontal scaling strategy Table of Contents
  1. HTTP vs WebSocket: The Hybrid Approach
  2. System Architecture: Three Layers
  3. Inbox Pattern: Guaranteed Delivery
  4. Message Routing Between Chat Servers
  5. Group Chat & Fan-Out Pattern
  6. Online Presence & Heartbeats
  7. Scaling to Millions of Users

01HTTP vs WebSocket: Choosing the Right Protocol

Your first instinct when designing a chat system might be to use HTTP for everything — it's familiar, it scales well, every load balancer knows how to handle it. And for sending a message, HTTP works perfectly: client makes a POST request, server receives the message and returns a 200 OK. Clean, simple, stateless. But receiving messages instantly is where HTTP fundamentally breaks down.

HTTP is client-initiated. The server cannot push a message to a client without being asked first. The two naive workarounds are short polling (client asks "any new messages?" every few seconds — wasteful, generates enormous load for empty responses) and long polling (client opens a request and the server holds it open until a message arrives or a timeout occurs — better, but each pending request still consumes a server connection). Neither is acceptable for a real-time chat system at scale.

WebSockets solve this properly. After an initial HTTP handshake (which upgrades the connection), a WebSocket is a persistent, full-duplex channel — the server can push messages to the client any time, without the client asking. Latency drops to milliseconds. But WebSocket connections are stateful: each connection ties up a server process or thread, scaling requires careful architecture, and load balancers need WebSocket awareness (sticky sessions or a shared state layer).

The industry-standard solution is a hybrid approach: HTTP for sending (stateless, scales trivially with horizontal load balancers), WebSocket for receiving (persistent connection managed by dedicated chat servers). This is exactly what WhatsApp, Facebook Messenger, and Discord use. You get HTTP's scaling simplicity for the high-volume write path and WebSocket's push capability for the latency-critical receive path. Don't try to over-engineer a single-protocol solution — the hybrid is the answer.

💡 Server-Sent Events: The Middle Ground

Here's the thing most tutorials miss: Server-Sent Events (SSE) is a third option between long polling and WebSockets. SSE is HTTP-based (works through standard HTTP/2 multiplexing), server-push only (client-to-server still uses regular HTTP requests), and has automatic reconnection built in. For read-heavy, server-push scenarios like notifications, activity feeds, or dashboard updates, SSE is simpler to scale than WebSockets because it works over standard HTTP/2 connections without the complexity of WebSocket upgrade management. The tradeoff: SSE is one-directional (server → client only), while WebSockets are full-duplex. For chat, use WebSockets. For notifications, SSE is worth considering.

websocket_chat.js — client connection
Hybrid approach: HTTP for sending, WebSocket for receiving

1. Connect to chat server via WebSocket (persistent receive channel)
const ws = new WebSocket('wss://chat.example.com/ws?token=' + authToken);

ws.onopen = () => {
 console.log('Connected to chat server');
};

ws.onmessage = (event) => {
 const msg = JSON.parse(event.data);
 displayMessage(msg);
 Send acknowledgement back to server
 ws.send(JSON.stringify({ type: 'ack', messageId: msg.id }));
};

Heartbeat: send ping every 30 seconds to keep connection alive
setInterval(() => {
 if (ws.readyState === WebSocket.OPEN) {
 ws.send(JSON.stringify({ type: 'ping' }));
 }
}, 30000);

2. Send message via HTTP POST (stateless, load-balanced separately)
async function sendMessage(recipientId, text) {
 const response = await fetch('/api/messages', {
 method: 'POST',
 headers: { 'Content-Type': 'application/json', Authorization: `Bearer ${authToken}` },
 body: JSON.stringify({ to: recipientId, text, clientMsgId: generateId() })
 });
 return response.json();
}

02System Architecture: Three Layers Working Together

A chat system at scale isn't one server — it's a collection of specialized services organized into layers. The clean architectural separation is what makes it possible to scale each component independently as different bottlenecks emerge. Imagine the entire system divided into three functional layers, each with a distinct responsibility.

The stateless services layer handles everything that doesn't require persistent connections: user authentication, profile management, friend lists, group membership, and REST API endpoints for sending messages. These services are classic horizontally-scalable web servers behind a load balancer. You can add or remove instances freely without coordination — no state, no problem. This is where most of your write traffic lands (sending messages), and it's where you get the simplest scaling story.

The stateful services layer — the chat servers — is the interesting part. Each client maintains a persistent WebSocket connection to exactly one chat server. Modern servers can hold 100,000 to 1,000,000 concurrent WebSocket connections depending on hardware and configuration. As your user base grows, you add more chat servers. The critical constraint: since users are pinned to specific chat servers, routing messages between servers requires coordination — which we'll address in section 4.

The third-party integration layer handles what happens when users are offline: push notifications via Apple's APNs (Apple Push Notification service) and Google's FCM (Firebase Cloud Messaging). These services maintain persistent connections to literally hundreds of millions of devices and know how to deliver notifications even when apps are not running. Reinventing this infrastructure would be a multi-year project. You leverage it instead.

⚡ Service Discovery Is the Hidden Critical Component

The stateful chat server layer requires a service discovery mechanism that the architecture overview often glosses over. Every time a user connects or disconnects, the system needs to update a registry mapping "user X is on chat server Y." When Alice sends a message to Bob, Alice's chat server must look up which server Bob is currently connected to. This registry (often backed by Redis or etcd) is queried millions of times per minute. Its availability and latency directly determines message delivery latency. Design it for high read throughput, geographic distribution near your chat server regions, and automatic cleanup when connections drop.

LLM next token prediction probability distribution diagram showing vocabulary tokens with probability bars and sampling mechanism

03The Inbox Pattern: Guaranteeing Delivery When Users Go Offline

Here's the fundamental problem: if you only deliver messages over live WebSocket connections, you lose messages whenever a user is offline, enters a tunnel, or their app crashes mid-delivery. You need a reliability mechanism that survives disconnections. The inbox pattern is the industry solution.

Every user has a personal message queue called an inbox. When a message is sent to you, the system first attempts immediate delivery via WebSocket. If you're online and connected, your device receives the message through the WebSocket and sends back an acknowledgement (ACK) — a confirmation that delivery succeeded. If the ACK is received, the message is never persisted to your inbox; the happy path is zero-storage. Only when you're offline or when delivery fails does the message get written to your inbox, where it waits until you reconnect.

The ACK mechanism is what provides the "at-least-once" delivery guarantee. The server only considers a message delivered when it receives the ACK. No ACK means the server will retry. When you come back online and reconnect to a chat server, your server immediately checks your inbox, downloads any waiting messages to your device in order, collects ACKs for each, and only removes messages from the inbox after successful ACK. This ensures that even a crash mid-download is recoverable — reconnect and the download resumes from the last un-ACKed message.

The design choice to store messages only temporarily (for offline delivery) rather than permanently is intentional and different from email. This keeps the inbox small, cheap to query, and avoids the indefinite storage costs of maintaining complete message history for billions of users. If persistent message history is a product requirement, it's a separate, explicitly designed storage system — not the inbox.

✅ Design for Idempotency: ACKs Can Be Lost Too

The ACK itself can fail. What if the server sends a message, the client receives it and sends an ACK, but the ACK gets lost? The server retries, the client receives a duplicate. Your client must deduplicate messages by message ID. Generate message IDs that are globally unique and deterministic (include sender ID, timestamp, and sequence number). The client-side message store should be keyed by message ID, making duplicate delivery a no-op rather than a duplicate display. This idempotency requirement applies to every message operation in the system — including the client's own sent messages, which should use client-generated IDs to prevent duplicates from network retries.

inbox_pattern.py — server-side delivery logic
from dataclasses import dataclass
from enum import Enum
import asyncio

class DeliveryStatus(Enum):
 DELIVERED = "delivered" # ACK received — don't store in inbox
 PENDING = "pending" # in inbox, waiting for ACK
 FAILED = "failed" # retries exhausted

async def deliver_message(server, msg, recipient_id: str):
 # Try immediate delivery via WebSocket first
 ws_connection = server.get_connection(recipient_id)

 if ws_connection and ws_connection.is_alive:
 try:
 await ws_connection.send_json(msg.to_dict())
 # Wait for ACK with timeout
 ack = await asyncio.wait_for(
 server.wait_for_ack(msg.id), timeout=5.0
 )
 return DeliveryStatus.DELIVERED # ACK received — done!
 except asyncio.TimeoutError:
 pass # fall through to inbox

 # Delivery failed or user offline — put in inbox
 await inbox.store(recipient_id, msg, ttl_hours=168) # 7 day TTL

 # Trigger push notification for offline users
 await push_service.notify(recipient_id, {
 "title": ff"New message from {msg.sender_name}",
 "body": msg.preview_text,
 "data": {"message_id": msg.id}
 })
 return DeliveryStatus.PENDING

async def sync_inbox_on_connect(server, user_id: str):
 """Called when user reconnects — deliver all queued messages"""
 messages = await inbox.get_pending(user_id)
 for msg in messages: # in order
 await ws.send_json(msg.to_dict())
 ack = await wait_for_ack(msg.id, timeout=10.0)
 if ack:
 await inbox.remove(user_id, msg.id) # only remove after ACK!

04Message Routing Between Chat Servers

With multiple chat servers each holding thousands of WebSocket connections, a new problem emerges: when Alice (connected to Chat Server 1) sends a message to Bob (connected to Chat Server 3), how does the message get to Bob's server? This is the inter-server routing problem, and it's where naive architectures fail at scale.

The solution uses a user presence service — a fast key-value store (Redis works well here) that maps each online user to their current chat server. When a user connects, their chat server registers them: presence.set("user:bob", "chat-server-3"). When a user disconnects, the registration is removed. When Alice's server needs to deliver a message to Bob, it first queries the presence service to find Bob's server, then makes a direct RPC (Remote Procedure Call) to that server to push the message through Bob's WebSocket connection. This direct server-to-server approach minimizes latency — the message takes one extra hop instead of going through a central broker.

The counterintuitive insight: a message queue (like Kafka) between chat servers seems like a natural fit but actually adds latency for the online delivery path. For real-time chat, you want the lowest possible latency — one direct RPC hop between servers. Kafka is excellent for the offline path (storing messages durably before inbox processing) but adds unnecessary latency to the critical hot path. Discord's architecture is a well-documented example of this direct server-to-server RPC approach.

⚠️ The Split-Brain Problem with Presence

What happens when the presence service reports Bob is on Chat Server 3, but Bob's connection actually dropped 2 seconds ago before the presence entry expired? Alice's server makes an RPC to Chat Server 3, which has no connection for Bob. This is the split-brain window — the gap between when a connection drops and when the presence service reflects it. The solution: Chat Server 3 should respond to the "no connection found" case by queuing the message to Bob's inbox immediately and returning success to Alice's server. Never fail silently when presence is stale — always fall through to the inbox as the reliability backstop.


05Group Chat & the Fan-Out Pattern

One-on-one messaging is straightforward: one sender, one recipient, one delivery path. Group chat with 100 members is fundamentally different: one message must be delivered to potentially 100 different users, each possibly connected to different chat servers, each possibly online or offline. The naive solution — a for-loop that runs the one-on-one delivery logic 100 times — actually works at this scale, but requires careful design to handle the fan-out efficiently.

The fan-out pattern works as follows: when a message is sent to a group, the server retrieves the list of group members. For each online member, it makes an RPC to their chat server to push the message through their WebSocket. For each offline member, the message is written to their inbox. The fan-out to 100 members can be parallelized — all RPC calls go out concurrently, not sequentially, keeping latency low. At 100 members, this is manageable. At 1,000 members (like Slack workspaces), the fan-out requires more careful batching and queuing. At 10,000+ members (broadcast channels), a different architecture — publish-subscribe with per-channel queues — is necessary.

Here's the thing most system design explanations miss: group chat requires sender-side deduplication. If Alice sends to a group, Alice's own device should receive the message back (so she sees it in the conversation thread) — but only once. The server must include the sender in the fan-out with a flag indicating "this is your own message" so the client handles it correctly, and use the client-generated message ID to prevent the client from displaying a duplicate alongside the locally-rendered optimistic message.

🔬 Fan-Out at WhatsApp Scale Is a Distributed Systems Problem

For groups with millions of members (broadcast channels, community groups), the fan-out pattern must shift from write-time fan-out to read-time fan-out. Instead of writing the message to each member's inbox at send time, you write the message once to a channel queue and let each member's client "pull" it when they reconnect. This is the distinction between "fan-out on write" (hot path, high write amplification) and "fan-out on read" (lazy, lower write cost but more complex read logic). Most chat systems use write-time fan-out for small groups (< 100 members) and read-time fan-out for large channels (> 1000 members).


06Online Presence: The Green Dot Problem

The green "online" indicator seems trivial — just check if the user has an active WebSocket connection, right? The problem is that WebSocket connection state is an unreliable indicator of user presence. A phone screen turns off, the app goes to the background, the phone enters a tunnel — the WebSocket connection can appear alive at the TCP level for minutes after the user has effectively gone away. You don't want to show your friend as "online" when they're actually asleep and their phone is in airplane mode.

The solution is application-level heartbeats. Clients send a WebSocket ping frame (or a small JSON message) to their chat server every 30 seconds. The chat server tracks the timestamp of the last received ping for each connected user. If no ping has been received within 60 seconds, the server marks the user as offline — even if the underlying TCP connection is still technically open. This 60-second grace period smooths out brief disconnections: a 10-second tunnel or elevator doesn't show your friend as offline and back online repeatedly, which would be annoying and consume bandwidth updating all their contacts.

Presence updates need to be broadcast to relevant contacts efficiently. When Alice goes online, her contacts need their presence indicators updated. The fan-out here is the same problem as message fan-out: efficiently propagating state changes to a potentially large number of subscribers. For most chat apps, presence updates use eventual consistency — there's no guarantee that all of Alice's contacts see the green dot at exactly the same millisecond, and that's fine. The 30-second heartbeat interval already means presence has 30 second granularity by design.

💡 Presence Updates Are High-Volume: Batch Them

For a user with 500 contacts, going online triggers 500 presence updates. Multiplied by millions of users doing this every time they open the app, presence becomes one of the highest-volume write operations in the system. The optimization: batch presence updates and use pub-sub channels rather than individual writes. Each user subscribes to a "contacts" presence channel. When Alice goes online, a single event is published to that channel, and a fan-out service distributes it. Rate limit and debounce presence broadcasts — don't broadcast "online" status more than once per 30 seconds per user, and batch updates to avoid a thundering herd when millions of users wake up simultaneously after an outage.


07Scaling from Thousands to Millions of Users

The architecture above works at 10,000 users. At 10 million users, several components will break under load, and you need to anticipate where. The first bottleneck is almost always chat server connection capacity. Each chat server handles a finite number of concurrent WebSocket connections. Scale this by adding more chat servers, updating service discovery to route new connections to servers with capacity. The bottleneck is linear and predictable — the easiest scaling problem in the system.

The inbox database is the second bottleneck. As millions of users go offline and come back online, their inboxes are written to and emptied constantly. The access pattern is heavily user-partitioned: user A's inbox is completely independent of user B's inbox. This makes horizontal partitioning (sharding) by user ID natural and effective. Shard the inbox database by user ID hash. Each shard handles a fraction of users. Add shards as load grows. This scales linearly and avoids cross-shard coordination.

For global reach, deploy chat servers in multiple geographic regions — US East, Europe, Asia Pacific, etc. Users connect to their nearest region for minimum latency. But cross-region messaging — Alice in London sending to Bob in Singapore — adds complexity: the message must travel from the London chat server cluster to the Singapore cluster. This is solved with a global message routing layer (often a Kafka-based backbone between regions) that handles the slower, higher-latency cross-region delivery path while keeping intra-region delivery fast and direct.

⚡ Connection Limits: Linux Kernel Tuning

A default Linux server configuration limits far fewer than 1 million concurrent connections. The default file descriptor limit (65,536 on most systems) caps the number of open sockets. To run high-concurrency WebSocket servers, tune: ulimit -n 1000000 (increase file descriptor limit), adjust net.core.somaxconn and net.ipv4.tcp_max_syn_backlog in kernel settings, and consider a socket-sharing architecture where multiple worker processes share the same listening port. With proper tuning, a single modern server (64+ cores, 256GB RAM) can sustain 500,000 to 1,000,000 concurrent WebSocket connections.


synthesisHow It All Connects: The Complete Message Journey

Follow a single message from Alice to Bob: Alice types "Hello" and hits send. Her client makes an HTTP POST to the stateless API servers (load-balanced, stateless). The API server receives the message, looks up Bob's presence in the presence service. Bob is online on Chat Server 3. The API server makes an RPC to Chat Server 3. Chat Server 3 pushes the message through Bob's WebSocket connection. Bob's client receives it, displays it, and sends an ACK. Chat Server 3 receives the ACK and confirms delivery. Total time: under 100 milliseconds for users in the same region.

If Bob is offline: the API server writes to Bob's inbox in the sharded inbox database and triggers a push notification via FCM/APNs. When Bob comes back online, he connects to a chat server (probably not Chat Server 3 — that's fine), the server checks his inbox, syncs all pending messages, collects ACKs, and cleans up the inbox. Bob sees the message immediately upon opening the app.


FAQFrequently Asked Questions

Why use WebSocket instead of HTTP for chat? + HTTP is client-initiated — the server cannot push messages to the client without the client requesting them first. For receiving messages instantly, HTTP forces you into polling (client repeatedly asks "any new messages?") or long polling (client holds a request open until a message arrives or timeout). Both are wasteful and slow compared to WebSockets. A WebSocket is a persistent, full-duplex channel established after an HTTP upgrade handshake. Once open, the server can push messages to the client at any time with millisecond latency. The hybrid approach — HTTP for sending (stateless, easy to scale), WebSocket for receiving (persistent, server-push) — gives you the best of both worlds: HTTP's load balancing simplicity for writes and WebSocket's real-time push for reads. What is the inbox pattern in chat system design? + The inbox pattern gives every user a personal message queue (inbox) that stores messages waiting to be delivered. When you're online, messages are delivered directly via WebSocket and never stored — the happy path is zero-storage. When you're offline or delivery fails, messages are written to your inbox and kept there until you reconnect. The critical reliability mechanism is acknowledgements (ACKs): when your device receives a message from the inbox, it sends an ACK. The server only removes the message from the inbox after receiving the ACK. No ACK means the message is retried. This creates an at-least-once delivery guarantee that survives disconnections, crashes, and network failures. Messages are typically stored in the inbox for a limited time (hours to days) rather than permanently. How does message routing work between chat servers? + Each user maintains a WebSocket connection to exactly one chat server. A user presence service (typically Redis) maps each online user to their current chat server. When Alice (on Server 1) sends to Bob (on Server 3), Alice's server queries the presence service to find Bob's server, then makes a direct RPC call to Bob's server, which pushes the message through Bob's WebSocket connection. This direct server-to-server routing minimizes latency. If Bob is offline (not in the presence service), Alice's server goes directly to the inbox path and triggers a push notification. The presence service requires careful handling of stale entries — when a connection drops, the server must update or expire the presence record to prevent routing to dead connections. What is the fan-out pattern for group chat? + Fan-out is the process of delivering one message to multiple recipients. In group chat, one message must reach every group member. The write-time fan-out pattern: when a message is sent, the server retrieves all group members, makes parallel RPC calls to deliver to all online members' chat servers, and writes to the inbox of all offline members. This can be parallelized — all RPCs go out concurrently — keeping latency low. Write-time fan-out works well for groups up to ~100 members. For large groups (thousands of members), the write amplification becomes too expensive. Large-group systems often use read-time fan-out: store the message once in a shared channel queue, and each member's client pulls it when needed. Most systems use write-time for small groups and read-time for large channels. How does online presence detection work in chat systems? + Online presence can't rely solely on whether a WebSocket connection is technically open, because connections can appear alive at the TCP level even when the app is backgrounded or the device is asleep. The standard approach is application-level heartbeats: clients send a ping frame (or small JSON message) to their chat server every 30 seconds. The server records the timestamp of the last received ping for each user. If no ping is received within 60 seconds, the server marks the user as offline — even if the TCP connection persists. The 60-second grace period prevents brief disconnections (tunnels, elevators, signal drops) from toggling the online indicator. Presence changes are broadcast to contacts via a pub-sub mechanism, typically with eventual consistency rather than strict real-time guarantees. How do you scale chat servers for millions of users? + Chat server scaling is linear and predictable: each server has a maximum concurrent WebSocket connection capacity (100,000–1,000,000 depending on hardware and configuration), and you add more servers as connection count grows. Service discovery automatically routes new connections to servers with available capacity. The inbox database is sharded by user ID — each database shard handles a fraction of users independently, with no cross-shard coordination needed. For global reach, deploy server clusters in multiple geographic regions and route users to their nearest region. Cross-region messages travel over a global routing layer (often Kafka-based). OS-level tuning (file descriptor limits, TCP buffer sizes) is required to unlock high connection counts on each server. What's the difference between APNS and FCM in chat apps? + Apple Push Notification Service (APNs) and Google Firebase Cloud Messaging (FCM) are third-party services that deliver push notifications to offline mobile devices. Both maintain persistent connections to hundreds of millions of devices and know how to wake up apps even when they're not running. APNs is for iOS/macOS devices; FCM covers Android and can also reach iOS. Your chat server sends a notification payload to APNs or FCM (via their APIs), specifying the device token and notification content. The platform service then delivers it to the device. This is infrastructure you absolutely don't want to build yourself — maintaining persistent connections to hundreds of millions of devices is a massive engineering undertaking. The tradeoff: you depend on Apple's and Google's reliability for offline delivery. How do you handle message ordering in a distributed chat system? + Message ordering in distributed systems is non-trivial because messages travel different paths (direct delivery vs inbox), servers have slightly different system clocks, and network delays vary. Three common approaches: Client-side timestamps (simple but unreliable — clients can have wrong clocks or manipulate them). Server-side sequence numbers (each message in a conversation gets a monotonically increasing sequence number assigned by the server — reliable but requires a central counter per conversation, which can become a bottleneck). Hybrid logical clocks or vector clocks (more complex, used in systems like CRDTs for collaborative editing). Most chat apps use server-side sequence numbers per conversation, which provides reliable ordering within a conversation and scales reasonably with per-conversation sharding.

💬 Chat System Lab

Four experiments: protocol comparison, inbox delivery, fan-out calculator, and server capacity planner.

Message delivery latency comparison across protocols

Protocol Latency Simulator Base network latency (ms) 50 Poll interval — short polling (s) 5s — Short polling (avg) — Long polling — WebSocket — Winner

Message lifecycle — click buttons to simulate scenarios

Inbox Pattern Simulator 0 Delivered (WS) 0 In inbox 0 Push notifications 0 Retries

Fan-out cost at different group sizes

Fan-Out Calculator Group members 50 % Online 60% RPC latency (ms) 5ms — Online (direct RPC) — Offline (inbox) — Fan-out latency — Strategy

WebSocket servers needed as DAU grows

Chat Server Capacity Planner Daily Active Users (DAU) 1M Concurrent % online 10% WS connections per server 100K — Concurrent users — Chat servers needed — Est. server cost/mo — Headroom
Tags
chat-system-designWebSocketinbox-patternfan-outpresenceAPNs-FCMhorizontal-scalingmessage-routingsystem-design-interviewreal-time-messaging
Share this article