System Design Ch.7: Chat System for 50,000 Users

Outline

Transcript

0:00 Picture this. It's 9:14 p.m. on a Sunday night. Okay, setting the scene. Yeah, right. So a breaking news incident suddenly drops, and instantly, like, 50,000 users flood into the exact same chat room. Oh, wow. Yeah, that's a lot of traffic all at once. Right. And someone in there types a tiny single-sentence message. And the moment that message hits the enter key, I mean, what actually happens? Absolute chaos if you're not prepared. Exactly, because if you're an engineer looking at this from, you know, a standard web architecture perspective, you might think, oh, it's just one simple database, right?

0:39 Right, like a standard API call, because that one tiny payload is essentially like a bomb going off in your back end. A bomb. Yeah, I mean, it triggers this absolutely massive wave of connection pressure, ordering pressure, presence churn, and just a giant delivery problem all at once. So it's not just logging a string of text. No, not at all. Standard request response infrastructure will just buckle under that within seconds. And honestly, that realization is the entire mission of this deep dive today.

1:06 We've got a stack of incredible architectural breakdowns here. And the core takeaway is that a chat system is definitively not just a standard CRUD application with faster writes. Right, because think about it. A CRUD app, it handles a request, does the work, sends a response, and then immediately forgets the caller. It's totally stateless at the connection layer. Exactly. But a chat product is fundamentally a live connection system, right, that just, you know, happens to also store messages. Yeah, the system actively has to remember the caller at all times.

1:39 Because maintaining that live state is literally the only mechanism that makes instant bidirectional delivery possible. Right. So for everyone listening, and we know you already know your databases, your queues, your caches, we are going to tear apart the mechanics of that live state today. Yeah, the real messy stuff. Exactly. We're answering four vital questions for you. First, who holds the millions of simultaneous connections? Second, how does the system calculate presence without, you know, melting the servers?

2:08 A very real risk, by the way. Oh, for sure. Third, how do we assign an unbreakable order to the messages? And finally, how does one single message fan out into delivery work for thousands of people? So let's start with the sheer physics of the problem. Because before the system can even process that 9:14 p.m. breaking news message, it first has to hold those 50,000 users. Well, actually, at a global scale, it has to hold millions of users simultaneously. Millions, right. Right. And if I'm building this, my first instinct is just to throw my standard API load balancers at it and, you know, let them route HTTP traffic.

2:43 Which is exactly how you end up taking down your own product. Wait, really? Why? Because, historically, systems tried to mimic live connections over standard HTTP using long polling. Sounds exhausting for the server. It is. The overhead of constant connection churn, rebuilding TCP sessions, parsing HTTP headers, and, you know, generating endless empty responses. It will absolutely destroy your CPU. And exhaust your server's ephemeral ports, I imagine. Exactly. Which brings us to WebSockets. With a WebSocket, you open one connection, and it just stays open for hours.

3:18 The client can push events up, and the server can push new messages down. Yeah. Typing indicators, read receipts, all immediately. But here's the catch. Holding millions of persistent sockets open requires a dedicated architectural layer. You need a connection gateway. Okay. Every open socket consumes memory and a file descriptor. And most users in a chat room, they aren't typing constantly. Right. They're mostly just lurking. Or the app is sitting open in the background. Exactly. So if you force your core business logic servers to hold hundreds of thousands of idle TCP connections, they're going to run out of memory just managing the keep-alives.

3:54 Oh, wow. I see. And worse, if a burst of heavy message processing suddenly spikes the CPU, that core server might freeze up for a second. And if it freezes, it misses the heartbeat pings from the clients. Yep. Dropping tens of thousands of users offline accidentally. Okay. That is the nightmare scenario right there. It really is. So connection gateways solve this by splitting the responsibility. The gateways are these incredibly lightweight services whose entire job is just to hold enormous numbers of mostly idle sockets efficiently.

4:25 Okay. So they handle the edge state. Right. They authenticate the user during the handshake. They track the heartbeats. And they terminate dead connections. They maintain a simple, localized map of which user-device pair is attached to which specific gateway process. Okay. So the gateway has my socket open and stable. But, you know, a socket isn't the feature users actually care about. The feature users actually care about is that little green dot. Oh, the presence indicator. Yeah. They want to know who else is online to discuss this breaking news.

4:53 So how does the gateway's raw TCP socket translate into a visual presence indicator for everyone else? Well, we have to shatter a major illusion here. Uh-uh. That online dot you see on your screen. It is a fabricated approximation. It is not an exact truth. What do you mean? When you drive into a subway tunnel, your phone does not politely send a TCP FIN packet to the server to say goodbye. Oh, right. It just vanishes. Yeah. The connection just drops silently into a black hole. The gateway only realizes you are gone because it stops receiving your periodic heartbeat pings.

5:28 Okay. So the back end is essentially just guessing you're still there until you miss too many check-ins. Pretty much. The mechanics usually involve a fast in-memory store like Redis. Okay. Redis. Yeah. When a gateway receives a heartbeat from your phone, it updates a key in Redis saying you are online and it attaches a TTL, a time to live, to that key. Like what? 30 seconds? Sure. Say 30 seconds. If your network drops, your heartbeats stop. But your Redis key doesn't disappear instantly. It just sits there, expiring slowly during that 30-second grace period.

5:59 And while I initially called that a bad user experience, I'm realizing that grace period is actually saving the entire back end from collapsing. It is the single most important defense mechanism for presence. I mean, if a system claims presence is absolute, real-time truth, it will spend all of its compute power fighting false transitions. Presence should be cheap to calculate and smoothed out. A small, built-in lag is vastly superior to a system that constantly flickers and triggers millions of useless update events.

6:29 Presence is a directional signal, not a guarantee. Okay, let's move from the green dot back to the actual data. The user has their connection, they see people are online, and they hit send on that 9:14 p.m. breaking news message. A fun part. Yeah. We now have to order it. How does the system prevent the humiliating bug where someone else's reply shows up in the timeline, like before the original message? Well, the foundational rule of distributed system design applies here. Never trust device clocks.

6:58 Never. Client timestamps are metadata. They are not truth. Because phones experience massive clock drift, right? Massive. Yeah. Or, like we mentioned earlier, a user might be in a subway. They type a message, hit send, and their phone buffers it for 10 minutes until they hit street level and get a signal. Right, and if the server trusts the timestamp on the phone, it will insert that message 10 minutes in the past, completely scrambling the flow of the conversation. Exactly. That forces the server to be the ultimate arbiter of time.

7:27 Production chat systems avoid that global complexity. The guarantee you actually need is heavily bounded. One room, one accepted order. Okay, so walk me through the strict right path to achieve that one room order. Sure. So the sender's gateway receives the raw payload over the WebSocket. The gateway looks at the routing map and forwards it to the specific internal chat service that owns that specific room. The room owner. Yeah. That owning service acts as the single authority. It does the validation, checks rate limits, ensures the user isn't banned, and then it appends the message to a durable conversation log in the database.

8:01 It has to hit disk. Absolutely crucial. Because only once it's safely in that durable log does the database assign it a monotonically increasing sequence number. Like sequence number 8,451. Exactly. And after the sequence number is assigned, only then does the server send the sent acknowledgement back to the user's phone. Durability always precedes delivery. Always. Once committed, that sequence number becomes the unchangeable law of the room. Every client fetching history just asks the server, hey, give me everything after sequence 8,451.

8:38 That monotonic sequence number is the anchor for the next massive hurdle, though. Which is getting the message out to everyone else. So the message is durably stored. It holds deli ticket number 8,451. Now, we have to fan it out to the 50,000 people staring at their screens. And this is where the architecture has to branch, right? Oh, definitely. Because the way you handle a private one-on-one chat will destroy your system if you apply it to a massive room. Scale forces a complete paradigm shift here.

9:07 So let's look at the small scale first. In private chats or small groups, the industry standard is fan out on write. Fan out on write. Yeah. The room owner writes the message to the database once. The server immediately looks up the connection gateways for the three other people in the group and actively pushes the delivery work out to them over their open WebSockets. And if they are offline, the server proactively updates their unread badge counters in the database. Right. It's a calculated bargain.

9:34 Paying for three delivery operations and three database updates for one new message is cheap. The newest messages are pre-calculated and waiting exactly where the clients expect them. It's like texting a few friends individually. You just push the data directly to their devices. But if you try to do that in the 50,000-person breaking newsroom... You trigger a cascading failure. That is the fan out amplification problem. I mean, imagine turning one single database write into 50,000 immediate routing decisions, 50,000 individual database updates for unread counters, and 50,000 push notifications.

10:08 Oh, man. The message queues will instantly back up. The CPU on the database will spike to 100%. One celebrity saying hello in a massive room would literally melt the infrastructure. So how do you actually stop the write amplification? At massive scale, the system abandons the push-heavy private messaging model and shifts to a pull-oriented log model. A pull model? Yeah. Giant rooms rely on a hybrid architecture. The message is written once to the central room timeline. It becomes a shared fact. And the clients are responsible for coming to get it.

10:40 Basically, yeah. Connected active clients might receive a tiny, lightweight, real-time event over their WebSocket just saying, hey, new message available. But not the whole message. No. Yeah. The system explicitly does not copy the full message payload to 50,000 individual inboxes. And it definitely doesn't update 50,000 unread counters in the database. That would be too much. Way too much. The per-user state shrinks dramatically down to just a set of lightweight cursors. Cursors, meaning like pointers to that deli ticket sequence number.

11:12 Yes. The database just tracks minimal metadata, like user A's last delivered sequence was 8,400, and user A's last read sequence was 8,390. Okay. So when user A opens their app, their local client looks at its own database, sees a gap, and fetches the missing sequence range. Okay. Let's stress test this. Let's go back to our user who was trapped in the subway train when the breaking news dropped at 9:14 p.m. They were completely disconnected. Ten minutes later, they walk up to street level. Their phone regains a cellular signal, and the WebSocket reconnects.

11:50 What is the actual sequence of events? Well, a naive architecture would have tried to keep all the messages they missed trapped in a massive queue in the gateway's memory, desperately waiting for them to reconnect. Which means during a major cellular provider outage, your gateways would try to queue messages for millions of offline users and run out of RAM in minutes. Precisely why we rely on the durable room log. When that mobile client reconnects, it doesn't ask the gateway what it missed. It looks at its local cursor.

12:19 It says to the server, hi, I'm back. Give me everything after sequence 8,451. And this approach is really the only thing that makes multi-device sync mathematically sane. Oh, true. If you have your phone, your laptop, and your tablet, they don't need to communicate with each other. Each device simply maintains its own cursor. The server treats every reconnecting device as a simple replay problem against the central log. Production systems rely on at-least-once delivery. Combined with idempotent deduplication on the client side.

12:52 Idempotent. Meaning the operation can be applied multiple times without changing the end result. But how does that actually work in practice on the phone? It's beautifully simple. The server might send sequence 8,452 three separate times because the network keeps dropping the acknowledgement packets. But the mobile client isn't running like a complex text-diffing algorithm to spot duplicates. It just attempts to insert the payload into its local SQLite database. And the sequence number is the primary key.

13:20 Boom. The SQLite database simply throws a unique constraint violation. It says, hey, I already have a row for sequence 8,452 and quietly drops the duplicate packet. Wow. At-least-once delivery guarantees the payload crosses the network. Idempotent design guarantees it doesn't ruin the user interface by showing the same message three times. That is clever. But what about retrieving ancient history? Like, if I'm catching up on a subway gap, I just need the last 10 minutes. But if I use the search bar to find a message from three years ago, does that hit the same database?

13:53 Absolutely not. That is why storage is heavily tiered into hot and cold paths. When you open a room, you want the last 50 messages instantly. The router looks at the requested sequence range. If it's recent, it hits a fast, append-heavy database or a ring buffer in memory. But a query for a three-year-old conversation has a completely different access pattern. Yeah. The router sees an ancient sequence number and directs that query to cold storage, usually a heavily compressed, column-oriented database that takes a few hundred extra milliseconds to respond, but costs a fraction of the price.

14:26 When you build this, you are actually building five different systems. Let's isolate those axes, because this is the critical takeaway for any engineer listening. Okay. Lay them out. Your connection gateways scale purely with the number of concurrent sockets. Your room ownership path scales with raw message throughput. Oh, yeah. Your presence infrastructure scales with heartbeat churn. Your push workers scale with offline fanout. And your history store scales with data retention. If you foolishly tie those all together in a monolith, a massive spike in message throughput like our 9:14 p.m.

14:57 breaking news event would starve the CPU. Which stops the server from responding to TCP keep-alives. Which drops everyone's sockets. Which forces 50,000 people to reconnect at once, taking down your entire infrastructure. Wow. But isolate the dimensions, and the system absorbs the shock gracefully. The right path might slow down by 50 milliseconds under extreme load, but the TCP connections stay open, presence keeps ticking, and the users never know the difference. Exactly. A beautifully designed chat system is ultimately a carefully bounded way to make stateful communication feel completely effortless.

15:34 The user just sees a green dot, a fast message, and everything in perfect order. Hiding that chaos is the true mark of great system design. It really is. And speaking of hiding chaos, we are going to tease what's coming up in Chapter 8 of this system design series. We are going to step away from live WebSockets and dive into the absolute nightmare of designing a massive notification system. Oh, that's a beast. It is. We'll be tackling multi-channel delivery across push, SMS, and email, managing strict user preferences, enforcing quiet hours, and dealing with distributed retries.

16:07 That's it for today's chapter. See you next time. Keep diving deep.