You type hello and it lands on your friend's phone before you lift your thumb. Behind that instant there is a connection that never hangs up, a server that remembers which socket is yours, and a receipt system that knows your friend's phone and laptop each saw something different. This lesson builds that system for 50 million daily users.
Outcomes
Start with the answers the interviewer in Ch. 12 gives when the candidate asks clarifying questions. They are numbered so later sections can cite exactly which one breaks.
(1) Clients: a mobile app and a web app.
(2) Modes: 1-on-1 chat and group chat, with at most 100 members per group.
(3) Features: low delivery latency, online presence, the same account on multiple devices at once, and push notifications. Text messages only.
(4) Message size: under 100,000 characters.
(5) Scale: 50 million daily active users.
(6) History: stored forever.
(7) Encryption: end-to-end encryption is not required for now.
Requirement (2) shapes the fanout math: every group message is a one-to-100 multiplication at most, which is cheap enough to copy per recipient. Requirement (5) shapes the connection math: millions of simultaneous sockets force a stateful tier that stateless API servers cannot provide. Requirement (6) shapes storage: forever means the history store must scale horizontally from day one.
Picture three ways to learn that a letter arrived. In the first, you walk to the mailbox every five minutes and usually find nothing. In the second, you walk over, and the clerk holds you at the counter until a letter arrives or the office closes, then you walk home and come back again. In the third, you and the clerk pick up a phone and simply stay on the line, either side talking whenever there is something to say.
That is the whole receiver problem in one image. The mailbox walk is polling: simple, but most trips return empty and each one costs a connection. Waiting at the counter is long polling: fewer empty answers, but the waiting still ends and restarts forever. The open phone line is WebSocket: one connection, speech in both directions, silence costs nothing.
On the sender side there is no problem to solve. The sender initiates, so plain HTTP works: open a connection, POST the message, and keep it alive across sends so TCP handshakes are not repeated. Ch. 12 notes that early Facebook chat sent messages exactly this way. The receiver side is where HTTP fails, because HTTP is client initiated and the server has news it cannot volunteer.
So Ch. 12 uses WebSocket for receiving and then, since the channel is bidirectional anyway, for sending too. One persistent connection per client, one mechanism to operate, and everything that is not real time (signup, login, profile) stays on ordinary HTTP request and response.
Follow one message from User A to User B. They sit on different chat servers, B might be offline, and B might own two devices. The flow below is the Ch. 12 happy path with the offline branch drawn alongside it.
How to read: Press play and walk the message from You. Watch it gain an ID, land in the queue, persist in storage, then take the socket path to User B when online or the push path when offline.
Your message reaches chat server 1 over your persistent WebSocket connection.
Chat server 1 obtains a message ID from the ID generator, which fixes the order of this message forever.
Chat server 1 drops the message into the message sync queue, the recipient's inbox.
The message is persisted in the key-value store, so history survives everything after this point.
If User B is online, the message is forwarded to chat server 2, which pushes it down B's open socket. If B is offline, push notification servers wake B's device instead.
Notice what the queue buys: chat server 1 returns after the enqueue, not after B's device answers, so a slow or absent recipient never stalls the sender. And notice the ID step coming before the queue: order is assigned once at intake, so retries and redelivery can never reorder the conversation.
The chat tier is the only stateful service in Ch. 12. Each client holds a persistent socket to one chat server and stays there while the server lives. Stateless API servers handle login and profile behind a load balancer, but the socket is sticky by nature: moving it means reconnecting. A component called service discovery (Apache Zookeeper in Ch. 12) keeps the registry and picks the best server for each login.
How to read: Press Add server and watch the counter. Four of sixteen connections move to ember and the rest stay put, because only the entries landing on the new server's arc get reassigned.
The memory math sets the scale: at roughly 10KB per connection, one million concurrent users need about 10GB just to hold sockets, and the book stresses this figure is rough and language dependent. One giant box could theoretically hold that, but a single box is a single point of failure, so the design spreads connections and lets service discovery reassign them when a chat server dies.
Group chat reuses the inbox idea with a multiplier. When User A writes to a group with B and C, the message is copied into B's sync queue and C's queue: one inbox per recipient, each client draining only its own. At the (2) cap of 100 members that is at most 100 small copies, which Ch. 12 accepts gladly, noting WeChat uses the same approach up to 500 members. Past that size the copies stop being a rounding error, which section 7 prices out.
| Table | Primary key | Why this key |
|---|---|---|
| 1-on-1 message | message_id | one conversation orders by a single ID sequence |
| Group message | (channel_id, message_id) | every query in a group runs inside one channel, so the channel partitions |
inboxes = {} # user_id -> list of messages (their sync queue)
_next_id = {} # channel_id -> local sequence number
def send_to_group(channel_id, sender, members, text):
# Local sequence: ordering only matters inside one channel,
# so a per-channel counter is sufficient (Ch. 12, message ID).
_next_id[channel_id] = _next_id.get(channel_id, 0) + 1
msg = {"channel": channel_id, "id": _next_id[channel_id],
"from": sender, "text": text}
for m in members: # the 1-to-100 multiplication from (2)
if m != sender:
inboxes.setdefault(m, []).append(dict(msg))
return msg
def fetch_new(user_id, cur_max, channel_id):
# Each device keeps its own cur_max_message_id. News is
# anything in my inbox for this channel with a higher id.
return [m for m in inboxes.get(user_id, [])
if m["channel"] == channel_id and m["id"] > cur_max]Multi-device sync falls out of the same structure. Each device keeps its own cur_max_message_id, and news is defined exactly as the code says: my ID as recipient, plus a store ID above my watermark. The phone and the laptop simply hold different watermarks over the same inbox, so neither device can steal the other's messages. For IDs, Ch. 12 lists three options in increasing strength: MySQL auto increment (unavailable in most NoSQL stores), a global 64-bit generator in the Snowflake style, or the local per-channel counter used above, which is the easiest correct answer whenever ordering never needs to cross channels.
Presence is a separate tier of presence servers speaking WebSocket to clients and writing status into the KV store. Login writes {status: online, last_active_at}, logout writes offline. The interesting case is neither: the user whose train enters a tunnel and whose socket dies without saying goodbye.
The naive fix marks the user offline the moment the socket drops and online when it returns, which turns every tunnel into a flickering green dot. Ch. 12 answers with a heartbeat: the client pings every 5 seconds, and only silence longer than a window (30 seconds in the book's worked example) flips the status to offline. The numbers are an illustration of the logic, not a prescription: the window trades stale green dots against flicker.
Delivery of each change uses publish and subscribe: every friend pair holds a channel, and a status event publishes once per channel for each friend to consume. That is trivially cheap inside the (2) cap of 100. For a group of 100,000 it manufactures 100,000 events per flap, so Ch. 12 switches strategy: fetch online status when a user enters the group or refreshes the friend list, and never fan it out eagerly.
Ch. 12 closes with storage and a set of stretch topics. History goes to a key-value store for three stated reasons: easy horizontal scaling, very low access latency, and immunity to the long tail problem where huge relational indexes make random access expensive. The field evidence is named: Messenger on HBase, Discord on Cassandra. Generic data such as profiles, settings, and friend lists stays in relational databases with replication and sharding.
Polling or long polling for receipt? Answer: polling wastes connections on empty answers, long polling breaks across stateless servers, so WebSocket carries both directions.
Why is the chat tier stateful while login is stateless? Answer: the socket persists per client and sticks to its server, so only discovery plus reassignment can scale it.
What happens when a chat server crashes mid-conversation? Answer: service discovery issues each orphaned client a new server, history reloads from the KV store by watermark, and the queue redelivers.
Why keep a message queue between intake and delivery? Answer: the sender returns after the enqueue, so offline recipients and traffic spikes never stall the send path.
When does per-recipient copying stop working? Answer: when the member count leaves the small-group regime toward 100,000, both inbox copies and presence events go from rounding error to overload.