Short answers to what readers ask most about this topic.
01How do you design a real-time chat application like WhatsApp?
Hold one WebSocket per device on stateless gateways, and keep a Redis registry of which gateway holds which user. Store each message in the database first with a per-conversation sequence number, ack the sender, then fan out through Redis pub/sub to the gateways. Offline devices sync from their cursor on reconnect and get a push notification meanwhile.
02How do chat apps guarantee message order?
They do not trust clocks. The server assigns a sequence number per conversation inside the same transaction that stores the message, and clients apply messages strictly in that order. If a client sees a number that is not last plus one, it knows there is a gap and requests a sync.
03What happens to messages when the recipient is offline?
They are already stored in the database, so there is no separate offline queue. When the device reconnects it asks for everything newer than its cursor and receives it in pages. While it is offline, a push notification through FCM or APNs tells the user something arrived.
04Can Redis pub/sub be the only transport for chat messages?
No. Redis pub/sub is at-most-once, so a message published while a subscriber is disconnected is lost. Use it as a fast doorbell between servers while the database remains the source of truth that clients re-sync from.
05How do group chats scale without sending millions of messages?
Push directly to members for small groups, where deliveries equal members minus one. Above a threshold, publish one small new-message notice per gateway and let each client pull the content with the same sync query used on reconnect. A 100,000-member group pushed per member at one message a second would need about 100,000 deliveries per second by itself.
Design a Chat App Like WhatsApp: Real-Time System Design
How to design a real-time chat application like WhatsApp: WebSocket gateways, ack and ordering with per-conversation sequence numbers, offline delivery, presence and group fan-out.
To design a real-time chat app like WhatsApp, keep one WebSocket per device on stateless gateways, route messages through a Redis user-to-gateway registry and pub/sub, store each message first with a per-conversation sequence number, ack by that number, sync gaps on reconnect, and send a push notification when no connection exists.
In ERP and POS work, the realtime requirement is usually a badge that updates without a refresh, and if one update is missed the next page load repairs it. A chat application is the same plumbing with a much harder promise: a message must arrive once, in order, even if the phone was in a lift when it was sent.
I run Postgres, Redis and NestJS on a single VPS, so I have not operated millions of connections, and this post does not pretend otherwise. It designs the whole system from stated assumptions, shows every step of the arithmetic, and uses only parts I would actually reach for. The facts about WebSocket, Redis pub/sub, partitioning and push lifetimes come from the linked specifications and docs.
What does a WhatsApp-style chat system need to do, and how big is it?
Start with the behaviour you promise, because each promise costs something later. Four of them shape the whole design:
Delivery: every message reaches every member at least once, and the client removes duplicates so the user sees it once.
Order: messages in one conversation appear in the same order on every device. Order across different conversations does not matter.
Offline: a device that was off for a day catches up on reconnect, and a push notification wakes it when it is not connected.
Status: sent, delivered and read ticks, plus an online indicator, all cheap enough not to dominate the traffic.
Then size it. Every input below is an assumption I chose for the exercise, not a measurement; change the first block and the rest recomputes.
Assumed inputs (a design exercise, not a measurement):
daily active users = 10,000,000
messages sent per user/day = 40
peak vs average = 3x
online at peak = 20% of DAU
stored bytes per message = 200 (body + ids + timestamps)
share of group messages = 30%, average group size 20
connections per gateway = 50,000 (assumed; load-test it)
memory per connection = 10 KB (assumed; socket + buffers + session)
Messages
sent per day = 10,000,000 x 40 = 400,000,000
average sends / s = 400,000,000 / 86,400 = 4,629.6
peak sends / s = 4,629.6 x 3 = 13,889
storage per day = 400,000,000 x 200 B = 80 GB
storage per year = 80 GB x 365 = 29.2 TB
Connections
concurrent at peak = 10,000,000 x 20% = 2,000,000
gateways needed = 2,000,000 / 50,000 = 40
connection memory, total = 2,000,000 x 10 KB = 20 GB (500 MB per gateway)
Deliveries (one send fans out to every other member)
average recipients = 0.70 x 1 + 0.30 x 19 = 6.4
deliveries per day = 400,000,000 x 6.4 = 2.56 billion
average deliveries / s = 2,560,000,000 / 86,400 = 29,629.6
peak deliveries / s = 29,629.6 x 3 = 88,889
per gateway at peak = 88,889 / 40 = 2,222 / s
Per-recipient receipt rows (16 B each) would cost
2.56 billion x 16 B = 41 GB per day -> store a cursor per member instead
Three numbers drive the design. Peak sends of 13,889 per second is modest for a database. The 2,000,000 concurrent connections are the real constraint, and they are why gateways exist as their own tier. And the 88,889 peak deliveries per second show that fan-out, not storage writes, is the busy part: each message is stored once but delivered about 6.4 times.
How do WebSocket gateways connect users and route messages?
Each device holds one WebSocket to a gateway. RFC 6455 defines the persistent, full-duplex channel and the ping and pong control frames used to detect dead peers. The gateway is deliberately dumb: it authenticates the socket, forwards frames to the chat service, and writes frames to sockets. It holds no conversation state, so any gateway can be killed and clients simply reconnect to another.
// Gateway side. One entry per live device, refreshed while the socket is alive.
const ROUTE_TTL_S = 60; // twice the 30 s heartbeat, so one missed beat is survivable
async function onConnect(userId: string, deviceId: string, gwId: string) {
// Hash per user: fields are device ids, values are the gateway that holds the socket.
await redis.hset("route:" + userId, deviceId, gwId);
await redis.expire("route:" + userId, ROUTE_TTL_S);
}
async function onDisconnect(userId: string, deviceId: string) {
await redis.hdel("route:" + userId, deviceId);
}
// A stale field (gateway crashed, TTL not yet expired) only means one publish goes
// to a gateway that no longer has the socket. It drops the frame; the message is
// already in the database, so nothing is lost.
The one thing the system needs is a map from user to the gateway holding their socket. A Redis hash per user, with a short TTL refreshed by the heartbeat, is enough, because the map is a hint. The database is the truth.
Routing approach
Pub/sub messages at peak (derived)
What breaks
Broadcast every message to all gateways
13,889 per second x 40 gateways = 555,556 received per second
Each gateway filters out most of what it receives; cost grows with gateway count.
One channel per user
At most 88,889 per second, but 2,000,000 channels
Subscribe and unsubscribe churn on every connect and disconnect.
Route registry plus one channel per gateway
At most 88,889 per second, usually far fewer after batching per gateway
Needs the registry, and stale routes must be tolerated. This is the default I would start with.
The sender looks up all recipients in one pipelined round trip, groups them by gateway, and publishes once per gateway. A 20-member group spread over 3 gateways costs 3 publishes, not 19. The comparison above is why: broadcasting makes 555,556 receives per second, about 6 times the 88,889 ceiling of the targeted design.
// Every gateway subscribes to exactly one channel: its own id.
const sub = redis.duplicate();
await sub.subscribe("gw:" + MY_GATEWAY_ID);
sub.on("message", (_channel, raw) => {
const { userIds, frame } = JSON.parse(raw);
for (const uid of userIds) {
for (const socket of localSockets.get(uid) ?? []) socket.send(JSON.stringify(frame));
}
});
// Sender side: group recipients by gateway, then publish ONCE PER GATEWAY, not per user.
async function deliver(frame: MsgFrame, memberIds: string[]) {
const pipe = redis.pipeline();
memberIds.forEach((uid) => pipe.hgetall("route:" + uid));
const routes = await pipe.exec(); // one round trip for the whole group
const byGateway = new Map<string, string[]>();
const offline: string[] = [];
memberIds.forEach((uid, i) => {
const devices = Object.values((routes[i][1] as Record<string, string>) ?? {});
if (devices.length === 0) offline.push(uid);
for (const gw of new Set(devices)) {
byGateway.set(gw, [...(byGateway.get(gw) ?? []), uid]);
}
});
for (const [gw, userIds] of byGateway) {
await redis.publish("gw:" + gw, JSON.stringify({ userIds, frame }));
}
return offline; // these users get a push notification instead
}
Scaling out is then mostly arithmetic. At 50,000 connections per gateway you need 40 gateways for 2,000,000 sockets, and each handles about 2,222 deliveries per second. If you want the NestJS and Redis wiring in detail, my posts on WebSockets with NestJS and on scaling WebSockets with Redis pub/sub cover it; this post is about what sits around it.
Redis pub/sub is at-most-once: a message published while a subscriber is disconnected is gone for good. That is acceptable here only because pub/sub is a doorbell. The message is already committed in the database, and the client re-syncs from its cursor. Never use pub/sub as the only copy of a chat message.
How does a message flow from sender to recipient, and what is acked?
Give every message a client-generated id, called mid below. The client keeps retrying a send with the same mid until it sees an ack, and the server uses the mid to make retries harmless. The frames look like this:
client -> gateway {"t":"send","cid":"c_91","mid":"7f3a-41","body":"on my way"}
gateway -> client {"t":"ack","mid":"7f3a-41","cid":"c_91","seq":1042} // stored: first tick
gateway -> peers {"t":"msg","cid":"c_91","seq":1042,"from":"u_7","body":"on my way"}
peer -> gateway {"t":"recv","cid":"c_91","upTo":1042} // delivered: second tick
peer -> gateway {"t":"read","cid":"c_91","upTo":1042} // read
The ack means stored, not delivered. That is the first tick. The handler therefore does its work in a fixed order: dedupe, allocate the sequence number, insert, commit, ack the sender, and only then fan out.
async function handleSend(conn: Conn, f: SendFrame) {
const { seq } = await pg.tx(async (tx) => {
// Retry of a message we already stored? Return the SAME seq, do not insert again.
const dup = await tx.oneOrNone(
"SELECT seq FROM messages WHERE conversation_id = $1 AND sender_id = $2 AND client_msg_id = $3",
[f.cid, conn.userId, f.mid],
);
if (dup) return dup;
// The row lock serialises senders in THIS conversation only. If the INSERT below
// fails, the transaction rolls back and the counter goes back with it: no gap.
const { seq } = await tx.one(
"UPDATE conversations SET last_seq = last_seq + 1 WHERE id = $1 RETURNING last_seq AS seq",
[f.cid],
);
await tx.none(
"INSERT INTO messages (conversation_id, seq, sender_id, client_msg_id, body) VALUES ($1, $2, $3, $4, $5)",
[f.cid, seq, conn.userId, f.mid, f.body],
);
return { seq };
});
conn.send(JSON.stringify({ t: "ack", mid: f.mid, cid: f.cid, seq })); // after the commit
const offline = await deliver({ t: "msg", cid: f.cid, seq, from: conn.userId, body: f.body }, await memberIds(f.cid));
await pushFallback(offline, f.cid, seq);
}
Acking after the commit is the rule that matters. If you ack first and crash before the insert, the sender believes a message exists that never will. If you commit and crash before the ack, the sender retries with the same mid, the dedupe query finds the stored row, and the same seq is returned. At-least-once on the wire plus idempotent storage gives exactly-once to the reader.
How do you keep messages in order, and what does the storage schema look like?
Order by a number the server assigns per conversation, never by timestamps. Phone clocks drift and a message can be delayed in transit, so timestamps order messages by who has the worst clock. A per-conversation counter is simple, and it also gives clients a cheap way to detect a gap: if the next seq is not last plus one, something was missed.
-- Wrong: order by the sender's clock. Two phones 3 seconds apart disagree on what came first.
SELECT body FROM messages WHERE conversation_id = $1 ORDER BY client_sent_at;
-- Right: order by the number the server assigned inside the transaction.
SELECT body FROM messages WHERE conversation_id = $1 ORDER BY seq;
// Client rule: apply messages strictly in seq order per conversation.
function onMsg(m: Msg) {
const last = lastSeq.get(m.cid) ?? 0;
if (m.seq <= last) return; // duplicate, ignore
if (m.seq > last + 1) { requestSync(m.cid, last); return; } // gap: fetch last+1 onward
render(m);
lastSeq.set(m.cid, m.seq);
}
A global counter would put 13,889 writes per second through one row. The per-conversation counter row only serialises senders inside the same conversation. The schema below puts the counter on the conversation, makes the composite primary key the access path, and keeps delivery state as a cursor per member.
CREATE TABLE conversations (
id bigint PRIMARY KEY,
kind text NOT NULL CHECK (kind IN ('direct', 'group')),
last_seq bigint NOT NULL DEFAULT 0 -- the per-conversation counter
);
-- Delivery state is a cursor per member, not a row per message per recipient.
CREATE TABLE members (
conversation_id bigint NOT NULL REFERENCES conversations (id),
user_id bigint NOT NULL,
last_delivered_seq bigint NOT NULL DEFAULT 0,
last_read_seq bigint NOT NULL DEFAULT 0,
PRIMARY KEY (conversation_id, user_id)
);
CREATE INDEX members_user_idx ON members (user_id);
CREATE TABLE messages (
conversation_id bigint NOT NULL,
seq bigint NOT NULL,
sender_id bigint NOT NULL,
client_msg_id uuid NOT NULL, -- idempotency key from the client
body text NOT NULL,
created_at timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (conversation_id, seq),
UNIQUE (conversation_id, sender_id, client_msg_id)
) PARTITION BY HASH (conversation_id);
CREATE TABLE messages_p0 PARTITION OF messages
FOR VALUES WITH (MODULUS 16, REMAINDER 0); -- repeat for remainders 1 to 15
-- 29.2 TB per year / 16 partitions = about 1.8 TB each. Partitioning keeps indexes
-- small; it does not make one server hold 29.2 TB. Spread partitions across hosts.
Per-recipient receipt rows would be 2.56 billion rows or 41 GB per day by the arithmetic above. A cursor per member costs one row per membership and is overwritten in place. Partitioning by hash of conversation_id follows the Postgres declarative partitioning docs, and it keeps each conversation inside one partition so range scans on seq stay local.
A conversation that gets very hot is bounded by one row lock. If one group ever sends faster than a single transaction per round trip can carry, move the counter into Redis with INCR and accept that a crash can leave gaps. Clients already treat a gap as a reason to sync, so the design tolerates this trade.
How do you deliver messages to offline users, and when do you fall back to push?
Offline delivery is not a separate queue. The message is already in the database, so catching up is a query. On connect the client sends hello, and the server compares each membership cursor to the conversation counter and returns only what is newer, in bounded pages.
-- On reconnect the client sends only: {"t":"hello"}. The server owns the cursors.
SELECT m.conversation_id, m.last_delivered_seq, c.last_seq
FROM members m
JOIN conversations c ON c.id = m.conversation_id
WHERE m.user_id = $1
AND c.last_seq > m.last_delivered_seq; -- only conversations with something new
-- Then, per conversation, a bounded page (primary-key range scan, already ordered):
SELECT seq, sender_id, body, created_at
FROM messages
WHERE conversation_id = $1 AND seq > $2
ORDER BY seq
LIMIT 200;
When the deliver step finds no route for a recipient, it sends a push notification. Use the payload only to say there is something new, and put the conversation id and seq in the data fields. FCM lets you set a message lifetime of up to four weeks, and a collapse key means a burst of messages produces one pending notification rather than fifty. Web Push has the same idea: RFC 8030 requires a TTL header and lets the push service hold the message until the device returns.
Because the collapse key throws away all but the latest notification, the notification can never carry the message itself. It only tells the app to open a socket and sync. That also keeps message text out of the push provider, which is a privacy gain and not just a design convenience.
How does group message fan-out work?
Fan-out is where chat systems get expensive. You choose when to do the work: on write, pushing to every member as the message arrives, or on read, letting members pull when they next look. Most systems end up with a hybrid.
Strategy
Cost shape
Use it for
Fan-out on write (push to every member)
Deliveries per message equal members minus one; lowest latency
Direct chats and groups up to a few hundred members
Fan-out on read (members pull)
One write per message; cost paid by each reader on open
Very large groups and broadcast-style channels
Hybrid (push a notice, pull the content)
One tiny publish per gateway plus a sync query per reader
Everything above the push threshold
The threshold is not a taste decision, it falls out of the capacity numbers. The 500-member cut-off below is an assumed starting point to tune with load tests.
Strategy by group size (push threshold of 500 members is an assumed starting point):
group of 20, 1 message/s -> 19 deliveries/s push each fine
group of 500, 1 message/s -> 499 deliveries/s push each fine
group of 100,000, 1 message/s -> 99,999 deliveries/s push each IMPOSSIBLE:
more than the whole system's assumed peak of 88,889 deliveries/s
Large group: publish ONE tiny "new seq" notice per gateway (40 publishes, not 99,999),
and let each client pull the page it needs with the same sync query as a reconnect.
The same sync query that serves a reconnect serves a large group, so the hybrid costs no new code path. Order is unaffected because the seq still comes from the one conversation counter.
How does presence (online and last seen) scale?
Presence looks free and is not. If every connection refreshes a key every 30 seconds, presence becomes the busiest write path in the system:
Heartbeat every 30 s from every connection (assumed), refreshing the route key's TTL:
2,000,000 connections / 30 s = 66,667 refreshes per second
peak message sends = 13,889 per second
66,667 / 13,889 = 4.8x as many presence writes as message sends
Batch per gateway (50,000 connections, one pipeline flushed each second):
50,000 / 30 = 1,667 commands in ONE round trip, 40 gateways -> 40 round trips/s
Two changes fix it. Gateways batch refreshes into one pipeline a second, so Redis sees 40 round trips per second instead of 66,667. And readers pull presence for the chat that is open instead of every contact being notified of every change.
// Gateway: remember who beat since the last flush, send one pipeline per second.
const beaten = new Set<string>();
socket.on("pong", () => beaten.add(socket.userId)); // RFC 6455 ping/pong control frames
setInterval(async () => {
const pipe = redis.pipeline();
for (const uid of beaten) pipe.expire("route:" + uid, 60);
beaten.clear();
await pipe.exec();
}, 1000);
// Reader side: no push to every contact. Ask only for the chat the user has open.
// "online" = the route key exists. "last seen" = a column written on disconnect.
const online = (await redis.exists("route:" + peerId)) === 1;
Reuse the route key for presence: online means the key exists. Write last seen to a column when the socket closes, and accept that a crash shows a stale last seen until the TTL expires. Both are honest approximations, and users cannot tell.
What should you decide before building a chat system?
Before writing code, you should be able to answer each of these in one sentence. If you cannot, that is the part of the design still open.
What is the ack contract: stored, delivered or read, and does the sender retry with the same client id?
Where does the per-conversation sequence number come from, and what happens to it on rollback?
Is the database the only copy of a message, with pub/sub treated as a doorbell?
What is the push threshold for group size, and what does a client do on a gap?
How are stale routes handled when a gateway crashes?
What does presence cost per second at your peak connection count?
None of these needs exotic technology. A single Postgres primary and a Redis instance will carry a small product comfortably; the design above is what you grow into, and the numbers tell you when each piece starts to matter.
A chat system is a storage problem wrapped in a routing problem. Store first and number messages per conversation, let the database be the truth, and treat every real-time channel as a hint that can be lost. Do that, and gateways, push and presence become pieces you can scale one at a time.