On this page
Tracks

System Design — Chat System

Last reviewed 11 Sept 2026

Part of the system design series. See the framework and building blocks first if you haven’t.

1. Requirements

Functional

  • 1:1 messaging and group chat (channels up to tens of thousands of members for the Discord-style case).
  • Delivery and read receipts, typing indicators, online presence.
  • Message history, searchable, available when a user comes back online after being offline for hours or days.
  • Push notification when the recipient has no live connection.

Non-functional

  • Low end-to-end delivery latency — WhatsApp-scale systems target under a second between distant users.
  • Massive concurrent long-lived connections — hundreds of millions of simultaneously-open sockets across the fleet.
  • No message loss, no duplicate delivery on retry (idempotency), ordering preserved per-conversation (not globally).
  • Horizontally scalable connection layer — a single gateway node holding all sockets for all users doesn’t work past a few hundred thousand connections.

2. Where it sits / high-level architecture

flowchart LR
A[Sender client] -- WSS --> G1[Gateway node A]
G1 --> MS[Message service]
MS --> MDB[(Message store\nsharded by conversation_id)]
MS --> REG[(Connection registry\nuser_id -> gateway_node)]
MS -->|recipient on node B| G2[Gateway node B]
G2 -- WSS --> C[Recipient client]
MS -->|recipient has no\nlive connection| PUSH[Push provider\nAPNs / FCM]
G1 -.heartbeat.-> PRES[(Presence store, TTL)]
WebSocket gateway with a connection registry for cross-node routing
  • Clients hold a persistent WebSocket (or long-lived HTTP/2 stream) to one of many stateless gateway nodes behind a load balancer.
  • The connection registry (Redis or similar, user_id → node_id) is what lets a message for a user connected to a different gateway node get routed correctly — every write path checks it.
  • The message service persists first, then attempts live delivery, then falls back to push. Persistence-before-delivery is the core invariant that makes retries safe.

3. Storage and delivery trade-offs

ConcernChoiceWhy
Message storageWide-column store (Cassandra/ScyllaDB-family), partitioned by conversation_id, clustered by a time-sortable message idWrite-heavy, append-only, always read by conversation, never joined — a relational store’s strengths (joins, ACID transactions) aren’t needed and its write throughput ceiling is
Message idSnowflake/ULID-style time-sortable idGives correct ordering within a conversation for free, without a global sequence bottleneck
Fan-out for large groupsAsync, via a queue partitioned by group_id, not synchronous push per memberA 10,000-member Discord server sending a message can’t block on 10,000 synchronous socket writes on the request path
PresenceRedis key with short TTL, refreshed by socket heartbeatAbsence of a heartbeat = offline; cheap, self-expiring, no explicit “user went offline” event needed
Ordering guaranteePer-conversation onlyA global order across all conversations isn’t a real product requirement and would force all writes through one serialization point

4. Deep dive — delivery guarantees and exactly-once-feeling delivery

There’s no true exactly-once delivery over an unreliable network, so the goal is effectively-once processing via idempotency, same as the notification service and outbox patterns elsewhere in this series:

  1. Client generates a client-side message id before sending. If the send is retried (network blip, no ack received), the server dedupes on that id — a duplicate send never creates a duplicate message.
  2. Server persists the message first (durable write to the message store), then attempts live delivery over the socket, then the recipient ACKs.
  3. If no ACK arrives within a timeout, the message is re-sent over the existing socket, or queued for delivery on reconnect/push if the socket is gone.
  4. On reconnect, the client sends its last-seen message id per conversation (or a global watermark); the server replays everything after that — this is what recovers from a dropped gateway node without any message loss, because the socket layer is stateless and disposable but the message store is not.
  5. “Delivered” and “read” states are just further idempotent events appended against the message, not a separate subsystem.

5. What real systems do today

  • WhatsApp runs on an Erlang/OTP backend — chosen specifically for cheap, massively concurrent lightweight processes, one per connection, which is a natural fit for hundreds of millions of long-lived sockets. WhatsApp is commonly cited as handling on the order of 100 billion messages/day with sub-second delivery between distant users.
  • Discord migrated its message store from Cassandra to ScyllaDB to handle trillions of stored messages. Concretely: their Cassandra cluster had grown to 177 nodes with p99 read latency of 40-125ms; after migrating to ScyllaDB they run the same workload on 72 nodes (a ~60% node reduction) with p99 read latency down to ~15ms and write p99 around 5ms, at roughly 2x storage density per node. The move was driven by hot partitions and GC pauses under Cassandra’s JVM-based design — exactly the kind of “which store, and why” trade-off worth naming in an interview.
  • Facebook Messenger historically ran on heavily sharded MySQL clusters, prioritizing ACID guarantees per shard while scaling horizontally by user/conversation sharding, versus the wide-column-first designs above — evidence that there isn’t one universally correct storage choice, it depends on the consistency and query-pattern priorities of the team.
  • Large-scale group delivery (Slack channels, Discord servers) commonly uses a queue partitioned by group_id (Kafka-shaped) so per-member delivery work is spread across consumer workers instead of blocking the sender’s request.

6. Scaling & failure

  • Bottleneck: one gateway node can’t hold enough concurrent sockets → horizontally scale stateless gateway nodes behind a load balancer; the connection registry is what makes any node able to route to any user regardless of which node they’re actually connected to.
  • Bottleneck: connection registry (Redis) becomes a hot single point under connect/disconnect churn → shard the registry by user_id hash, same pattern as sharding a rate limiter’s key space.
  • Bottleneck: message store hot-partitions on a mega-popular conversation (public server, huge group) → this is Discord’s actual documented problem under Cassandra; solved by moving to a store with better hot-partition behavior (ScyllaDB’s shard-per-core architecture) rather than trying to manually re-shard a single conversation’s partition.
  • Bottleneck: storing every message forever → tier old messages to cheaper cold storage after an inactivity window; the hot store only needs to serve recent history fast.

What happens when a gateway node dies: all sockets on that node drop. Clients detect the disconnect and reconnect to any available gateway node (the load balancer picks one, no affinity required). On reconnect, the client provides its last-seen message id per active conversation; the message service replays anything it missed from the durable message store. Presence for those users flips to offline once their heartbeat TTL expires (seconds, not instant) and comes back once they reconnect and heartbeat again. No message is lost because the socket/gateway layer holds no state that isn’t reconstructible — the message store is the only source of truth, exactly mirroring the ride-hailing geo-index and news-feed Redis patterns elsewhere in this series: keep durable state in a real store, treat the fast/ephemeral layer as disposable.

Interview follow-ups

  • “How do you guarantee a message isn’t lost if the recipient’s gateway node crashes mid-delivery?” — Persist to the durable message store before attempting live delivery; reconnect replays anything undelivered from the store using the client’s last-seen id.
  • “How do you avoid duplicate messages when a client retries a send after a timeout?” — Client-generated idempotency key deduped server-side; the same key never creates two stored messages.
  • “Design fan-out for a 50,000-member Discord server.” — Never synchronous per-member push on the request path; publish once to a queue partitioned by group_id, consumer workers handle per-member delivery/push asynchronously.
  • “Why a wide-column store instead of a relational database for messages?” — Access pattern is pure append + range-scan-by-conversation, no joins; wide-column stores give much higher sustained write throughput and better horizontal partitioning for that specific pattern, at the cost of losing cross-table joins/transactions you don’t need here.
  • “How does presence work without a ‘user went offline’ event?” — Heartbeat over the socket refreshes a short-TTL key; absence of a heartbeat is offline. Simpler and self-healing versus explicit disconnect events, which can be missed on a hard crash.
  • “Why per-conversation ordering instead of a single global order?” — A global order forces every write through one serialization point/id source, which doesn’t scale, and no real chat product needs to know that a message in conversation A happened before one in unrelated conversation B.
  • “What’s the actual number that forced Discord off Cassandra, and what changed?” — Their cluster grew to 177 nodes with p99 reads of 40-125ms due to hot partitions and JVM GC pauses; moving to ScyllaDB cut that to 72 nodes at ~15ms p99 reads — cite it as a concrete “storage engine choice has a measurable operational cost” example.

Sources: How Discord Migrated Trillions of Messages from Cassandra to ScyllaDB — ScyllaDB · How Discord Migrated Trillions of Messages to ScyllaDB — The New Stack · How Discord Moved Trillions of Messages to ScyllaDB — Hello Interview · Design a Chat App System — systemdesign.one newsletter · Designing WhatsApp Messenger — GeeksforGeeks