High Level Design

Design a Chat System

End-to-end design of a real-time chat system for 50M DAU — covering WebSocket protocol choice, message fanout, multi-device sync, presence management, and media storage.

September 18, 2026·Updated September 18, 2026

Step 1: Requirements#

Clarifying Questions#

QuestionAnswer
What type of chat?Both 1-on-1 and group chat
Mobile or Web?Both
Scale?50 million daily active users (DAU)
Group member limit?Maximum of 100 people
Core features?1-on-1 chat, group chat, online indicator (text only)
Message size limit?Less than 100,000 characters
End-to-end encryption?Not required initially
Chat history retention?Forever
Media sharing?Images and files

Functional Requirements#

#Requirement
1Send and receive messages in real time (1-on-1 and group)
2Group chats with up to 100 members
3Message receipts — sent, delivered, read
4Image and media sharing
5Online / last seen presence indicators
6Push notifications for messages received while offline
7Multi-device support (phone + laptop in sync)

Non-Functional Requirements#

#RequirementDetail
1Low latencyMessage delivery under 500 ms end-to-end
2High availability99.99% uptime; chat servers can restart without message loss
3DurabilityMessages must never be lost once acknowledged
4Eventual consistencyRead receipts and presence can lag slightly
5Horizontal scalabilityMust handle 50M DAU; scale chat and presence servers independently

Estimations#

Assuming 50M DAU, each sending ~40 messages/day:

MetricCalculationResult
Messages/day50M × 402 billion
Peak QPS (3× average)2B / 86,400 × 3~70,000 msg/s
Avg message sizetext + metadata~200 bytes
Write bandwidth70,000 × 200 B~14 MB/s
Storage/day2B × 200 B~400 GB/day
Storage/year400 GB × 365~145 TB/year

Step 2: Entities#

EntityResponsibility
UserProfile, auth credentials, device tokens
MessageContent, sender, recipient or group, timestamp, message ID
ConversationLogical channel between two users (1-on-1)
GroupMetadata (name, owner) and member list
PresenceOnline status and last_active_at timestamp per user
MediaReference (URL, MIME type, size) to a file in object storage

Step 3: APIs#

Send a message#

POST /messages
ParamTypeDescription
conversationIdstring1-on-1 conversation ID or group ID
contentstringText body (or empty when sending media)
mediaUrlstring?CDN URL of an already-uploaded media file

Returns the created message with a server-assigned messageId.

Fetch message history#

GET /conversations/{conversationId}/messages?before={messageId}&limit={n}

Returns a page of messages before the given cursor, newest first.

Initiate a media upload#

POST /media/upload/init
ParamTypeDescription
mimeTypestringe.g. image/jpeg
sizenumberFile size in bytes

Returns a time-limited pre-signed URL. The client uploads directly to object storage and then sends the resulting CDN URL in a message.

Update presence status#

PUT /presence/status
Body: { "status": "online" | "offline" }

Get user presence#

GET /users/{userId}/presence

Returns { status, lastActiveAt }.


Step 4: Data Flow#

1-on-1 message send#

message-flow-a-to-b

StepActorAction
①User A → Chat Server 1Message frame sent over the existing WebSocket connection
②CS1 → CassandraAssign messageId; persist message to the key-value store
③CS1 → RedisCheck presence: is B online? Which Chat Server is B connected to?
④CS1 → Message QueuePublish message event addressed to B (Kafka / Redis Pub-Sub)
⑤Queue → Chat Server 2CS2 — the server holding B's socket — consumes the event
⑥CS2 → User BMessage pushed over B's open WebSocket connection
⑦User B → CS2Client sends delivery ACK
⑧CS2 → CS1 (via MQ)ACK routed back through the queue to CS1
⑨CS1 → User A✓✓ "Delivered" receipt shown on A's screen

If B is offline at step ③: skip steps ④–⑨ and instead enqueue a push notification to APNs (iOS) or FCM (Android).

Group message send#

  1. Sender sends to Chat Server 1, same as above.
  2. Message is persisted once with a (groupId, messageId) primary key.
  3. A copy is written to each member's sync queue (individual inbox).
  4. Each member fetches new messages from their inbox independently — their client tracks cur_max_message_id to avoid duplicates.

Step 5: High-Level Design#

Why WebSockets?#

HTTP is request-response; the server cannot push. Polling and long-polling simulate push but waste connections or add latency.

websockets

WebSockets start with an HTTP upgrade handshake and then become a full-duplex, persistent channel. Both sides can send at any time with minimal overhead. This is the industry choice (WhatsApp, Slack, Discord).

HTTP Long PollingWebSockets
DirectionUnidirectionalBidirectional
LatencyPoll intervalInstant
Server loadHigh (constant reconnects)Low (one persistent connection)

Architecture overview#

high-level-design

Stateless services (behind a load balancer): auth, user profiles, group management, notification dispatch. These are ordinary REST APIs.

Stateful services: Chat Servers each hold many long-lived WebSocket connections. Clients stay pinned to the same server for the duration of a session. Service discovery (Apache Zookeeper) routes new connections to the least-loaded server geographically close to the client.

Presence Servers: dedicated nodes that track heartbeat events and publish status changes via pub/sub.

Push Notification Service: forwards messages to APNs / FCM for offline users.

Object Storage + CDN: media files live in S3/GCS; the CDN serves reads, keeping bandwidth off the chat servers entirely.

API Server vs Chat Server — when does each handle a request?#

api-server-vs-chat-server The simplest mental model:

  • API Server = stateless HTTP — anything that doesn't need to be live. Client makes a request, gets a response, done.
  • Chat Server = stateful WebSocket — anything that needs to be pushed in real time. The connection stays open for the entire session.

The typical flow when a user opens the app:

  1. POST /login → API Server — authenticate, receive token.
  2. GET /contacts, GET /chats → API Server — load initial data.
  3. WebSocket upgrade → Chat Server — one persistent connection is established and stays open.

After the WebSocket is established, all real-time activity travels over it:

ActionWhere it goesWhy
Send / receive messageChat Server (WebSocket)Must be pushed instantly
Typing indicatorChat Server (WebSocket)Real-time, ephemeral
Online / offline presenceChat Server (WebSocket)Event-driven, not polled
Read receiptsChat Server (WebSocket)Real-time acknowledgement
Load older messages (scroll up)API Server (HTTP)Not real-time; paginated range scan
Upload mediaAPI Server (HTTP) → S3 directlyLarge payload, not time-critical
Create / update profileAPI Server (HTTP)Stateless write
Fetch contacts or group infoAPI Server (HTTP)Stateless read

Why can't you just use one server type for everything?

Chat servers are stateful — each one tracks which user IDs are connected to it. A message for User B must be routed to whichever server User B is currently connected to. That's a fundamentally different scaling problem from the stateless API tier, which can be load-balanced freely because no session affinity is needed.


Step 6: Deep Dive#

Service discovery#

service-discovery

All Chat Servers register with Apache Zookeeper on startup (host, port, current connection count). When a client connects, an API server queries Zookeeper and returns the best server address. The client then opens a WebSocket directly to that address. If the server restarts, the client reconnects and Zookeeper routes it to a healthy node.

Multi-device synchronisation#

message-sync

Each device independently tracks cur_max_message_id. On reconnect (or periodically), it pulls all messages with messageId > cur_max_message_id from the key-value store. Because message IDs are monotonically increasing within a conversation, this is a single range scan — no deduplication logic needed.

Online presence management#

Login

login

On WebSocket connection: presence server sets status to online and records last_active_at.

Logout

logout

Explicit logout calls the API; status flips to offline and timestamp is preserved.

Disconnection (the tricky case)

heartbeat

A raw TCP drop looks identical to a flaky network. Marking users offline immediately causes status to flicker. Instead, clients send a heartbeat every 30 seconds. The presence server marks a user offline only after missing several consecutive heartbeats — typically 60–90 seconds of silence.

Status fanout

status-fanout

Presence servers use a publish-subscribe model: when User A goes online, an event is published; every friend who has a subscription channel with A receives the update. For large groups (1000+ members) pushing to all subscribers becomes expensive — the optimisation is to deliver presence lazily: fetch status when the user opens a conversation rather than pushing proactively.

Media handling#

Never store media blobs in the chat database.

Upload flow:

  1. Client requests a pre-signed URL (POST /media/upload/init).
  2. Server validates the request and returns a time-limited write URL scoped to a single object.
  3. Client uploads directly to S3/GCS — chat servers are bypassed entirely.
  4. Client sends a chat message containing the CDN URL.
  5. Recipients fetch media via CDN edge nodes.

Storage optimisations:

  • Deduplication: hash each file (SHA-256); if the hash already exists, reuse the stored object (saves 20–30% for forwarded images).
  • Lifecycle policies: move media older than 6 months to cold storage (Glacier); auto-delete media in ephemeral chats.
  • Chunked upload: split large files into 5 MB parts, upload in parallel, resume on failure.

End-to-end encryption#

Messages are encrypted on the sender's device before they touch the network. Only the sender and recipient hold the decryption keys. Chat servers relay ciphertext — they cannot read message content. This is the model used by WhatsApp's Signal Protocol.

Performance & resilience#

Graceful degradation: under traffic spikes, disable non-critical features first (read receipts, typing indicators, presence updates) so core message delivery continues uninterrupted.

Rate limiting: cap at 100 messages/minute per user and 1,000 API calls/hour. Reject excess at the gateway — prevents spam without blocking legitimate traffic.


Step 7: DB Choices#

ComponentDatabaseWhy
Chat messagesCassandra / HBaseWrite-heavy, append-only, horizontal sharding by (conversationId, messageId), proven at WhatsApp / Facebook scale
User profiles, groupsPostgreSQLRelational structure, ACID transactions for membership changes
Online presence, session stateRedisIn-memory key-value with TTL; heartbeat writes need sub-millisecond latency
Media filesS3 / GCSObject storage built for large blobs, native CDN integration, independent scaling
Service registryApache ZookeeperDistributed coordination, watch notifications for server health

Step 8: Follow-Ups#

Scaling to millions of concurrent connections#

A single Chat Server can hold ~100K WebSocket connections. At 50M DAU with, say, 20% concurrently online, you need ~100 servers. Use consistent hashing to shard users across servers so the message router can find the right target without broadcasting.

Large groups (1,000+ members)#

Writing one inbox copy per member at 1,000+ members on every message is expensive. Alternatives: store the message once and have each client pull from a shared groupId cursor, or introduce a fanout service that processes writes asynchronously and batches delivery.

Video and audio calls#

WebSockets add unnecessary relay overhead for real-time media. Use WebRTC for peer-to-peer audio/video — it has its own signalling, STUN/TURN servers, and codecs optimised for low-latency streams. The chat server acts only as the signalling channel to initiate the call.

Full-text search over billions of messages requires a dedicated index. Write message content to Elasticsearch asynchronously (via Kafka consumer). Apply per-user access control at query time so users can only search their own history.

Compliance and data residency#

Regulated industries (finance, healthcare) require message retention policies and audit logs. Store encrypted message archives in a separate compliant store (e.g. AWS GovCloud), with a retention manager that enforces deletion schedules independently of the live chat database.


FAQs#

Why does logout go through the API Server and not the Chat Server?#

Logout is a stateless write, not a real-time event. When a user logs out, three things need to happen:

  • Auth token is invalidated in the DB
  • Device tokens are cleaned up (to stop push notifications)
  • Presence status is flipped to offline in Redis

These are durable, transactional operations — exactly what API servers are built for. The Chat Server only owns the WebSocket connection; it doesn't own auth tokens or persistent presence records. Mixing that responsibility in would couple stateful and stateless concerns on the same node.

Compare this to a network drop — that has no explicit signal, so it's the Chat Server that detects it via missing heartbeats and then notifies the Presence Server.

EventPathWhy
Explicit logoutClient → API Server → Presence ServerIntentional write; token invalidation needed
Network drop / crashChat Server → Presence Server (heartbeat timeout)No explicit signal; only the WS layer knows

Does the client talk to the Presence Server directly?#

No — the Presence Server is an internal service with no public endpoint. The client only has two channels open: HTTP to API servers and a WebSocket to a Chat Server. The Presence Server is always reached backend-to-backend:

EventPath
Login / explicit logoutClient → API Server → Presence Server
Heartbeat (staying online)Client → Chat Server → Presence Server
Network dropChat Server (timeout) → Presence Server
Friend reads your statusClient → API Server → Presence Server → Redis

What happens when a user gets disconnected?#

Unlike logout, a network drop sends no explicit signal — the server has to figure it out on its own via missed heartbeats.

disconnection-flow

StepWhat happens
① ②Client sends a heartbeat ping every 30s over the WebSocket. Chat Server forwards it to the Presence Server, which refreshes the Redis TTL (key set to expire in 60s).
③Network drops — no logout, no signal. Client is silently gone.
④ ⑤Chat Server notices missed heartbeats at t=60s and t=90s. After N consecutive misses (threshold: 2–3), it treats the connection as dead.
⑥ ⑦Chat Server notifies Presence Server → Redis: status = offline, last_active_at = now.
⑧Presence Server publishes an offline event via pub-sub — friends subscribed to User A's channel are notified.

Why not mark offline on the first missed heartbeat? A flaky network can drop one ping without the user actually leaving. Marking offline instantly causes presence to flicker. The threshold (2–3 missed beats ≈ 60–90s delay) absorbs transient drops while keeping "last seen" reasonably accurate — the same approach WhatsApp uses.

EventWho detects itHow
Explicit logoutAPI ServerClient sends PUT /presence/status offline
Network dropChat ServerCounts consecutive missed heartbeats
Server restartZookeeper watchHealth check fails; clients reconnect

What happens when a Chat Server crashes?#

Three things need to be handled: detection, client reconnection, and in-flight message safety.

① Detection — Zookeeper

Every Chat Server maintains a session with Zookeeper on startup. When it crashes, the session expires within seconds. Zookeeper fires a watch event → API Servers and the load balancer remove the dead node from the registry immediately.

② Client reconnection

All clients on that server get a WebSocket connection closed event. Their client-side retry logic kicks in:

WS closed → backoff timer → GET /chat-server (API Server)
           → Zookeeper returns a healthy server
           → new WebSocket connection established

The client reconnects to a different Chat Server — session affinity only lasts for the life of the connection. After reconnect, the client pulls messageId > cur_max_message_id to catch up on anything missed.

③ In-flight messages — the tricky part

This depends on where in the pipeline the crash happened:

Crash momentMessage fate
Before CS wrote to CassandraLost — client never got ACK, so it retries
After Cassandra write, before publishing to QueueSafe in DB; recipient pulls it on reconnect
After publishing to QueueSafe — Kafka holds it durably; recipient's new CS consumes it

The key protection: CS only sends ACK to the sender after both the DB write and the queue publish succeed. If the client never gets an ACK, it retries. This gives at-least-once delivery.

④ Presence

The crashed server stops sending heartbeats on behalf of its users. The Presence Server detects this the same way as a network drop — Redis TTL expires after 60–90s → users marked offline → friends notified via pub-sub.

Summary:

CS crashes
  → Zookeeper detects (seconds)
  → removes from registry
  → clients get WS disconnect → reconnect to healthy CS
  → unACKed messages retried by client
  → queued messages delivered normally
  → presence marked offline after heartbeat timeout

The Message Queue (Kafka) is the main durability shield — once a message hits the queue, a server crash cannot lose it.