System Design Ch.2: Replication and Consistency

Outline

Transcript

0:00 I'm going to start today with a pretty uncomfortable truth for any backend engineer listening to this right now. Oh, all right. Let's hear it. Every single time you read from a database replica, you might be reading entirely stale data. Yeah. That is the dark secret, isn't it? Right. Because most of us, we just look at our cloud console, we see that replication is turned on, and we assume we are completely safe. We love those green check marks in Terraform. Exactly. But simply knowing, like, my data is replicated, tells you absolutely nothing about whether your reads are actually consistent.

0:36 You might literally be serving your user's data from five minutes ago. And the funny thing is, your system's metrics would still probably call that a resounding success. Exactly. It's wild. We want to equip you, the senior engineer, with the structured reasoning to stop blindly trusting those database boxes. Because as soon as you trust a system with data you actually care about, you know, user profiles, financial ledgers, inventory counts, you are forced to copy it. Right. And how a system handles the disagreements between those copies, that basically defines its entire personality.

1:09 It really does. The moment you have more than one copy of a piece of data, your entire job just becomes managing the inevitable arguments between those copies. So let's start with the baseline physics here. Like, why do we even copy data in the first place? And why is it such a massive headache? Well, the easiest way to avoid data disagreement is to just keep everything on one massive server. Right? Which sounds great until it isn't. Exactly. Any senior engineer knows that is just a ticking time bomb.

1:36 Because meantime between failures is a statistical certainty. A single disk failure or a misconfigured kernel update and boom. Yeah. Or a power surge in one specific availability zone and everything is just gone. You absolutely have to replicate to survive hardware reality. It's basically a classic insurance policy. I always think of it like having backup copies of highly sensitive corporate documents in different regional offices. Okay. Yeah. I like that. Right. So if the New York office burns down, it's a tragedy, obviously, but the company survives because the Chicago office has the backups.

2:12 Right. But the catch is, what happens if an executive in New York and an executive in Chicago both try to redline their copy of the document at the exact same millisecond? Yeah. Which version reflects reality? Exactly. And that synchronization problem is exactly what the most common replication pattern leader follower attempts to solve by force. By force. I like that. Like post-Greskual streaming replication. Exactly. That's the quintessential example. You designate one node as the absolute authority, the leader.

2:41 The boss. Right. The leader takes all the writes. Yeah. Every single insert, update, or delete has to go through that one node. And the followers. The followers just get a continuous stream of those changes. They just apply them in a strict order. So you completely eliminate the conflict because there is literally only one door to write through. Yes. You can read from any of the followers. Yeah. Which is great because it scales your read capacity horizontally, but your write path is totally centralized.

3:10 So it's super predictable. It is wonderfully predictable. Up until the leader chokes. Oh, man. Yeah. Which brings us to the disaster scenario. Every on-call engineer dreads. Hardware fails. Yeah. Networks drop packets. Suddenly, your followers stop receiving that replication stream. And they can't ping the leader. Right. And the system has to ask itself this terrifying question. Is the leader actually dead, or is the network just having a temporary hiccup? And that ambiguity is where the danger lives, right?

3:39 It is the most dangerous state for a distributed database. Because to maintain availability, your automated failover scripts will try to promote a new leader. But what if the original leader wasn't actually dead? Exactly. What if it was just bogged down by a massive garbage collection pause? Or a localized network partition? I have actually seen this happen in production. The system promotes a follower. The application starts sending live traffic to the new leader. And then the old leader wakes up.

4:09 Yes. It wakes up and just starts taking writes again. Because it has no idea it was fired. It cheerfully accepts those rights. So now you have two nodes, both completely convinced they're the primary, both accepting different rights from different users. It's the two managers problem. Oh, how so? Exactly like two engineering directors who suddenly can't reach each other on Slack because of a network outage. Oh, I've been there. Right. They both think the other one is out sick or offline. So they both start issuing contradictory instructions to the team.

4:39 Right. One is yelling, deploy the hot fix to production. And the other is yelling, no, roll back the entire cluster. And the team is just caught in the middle of creating an absolute mess. Exactly. They both genuinely believe they possess the ultimate authority. Well, in database terms, we call that split-brain. Split brain? That just sounds terrifying. It is. And the terrifying thing about split brain is that it produces conflicting state that no automated process can cleanly resolve. Like, give me an example.

5:09 Okay. So imagine old leader A accepts a right changing an auto-incrementing primary key to 105 for a new user. Okay. And new leader B assigns that exact same primary key, 105, to a completely different user. Oh, no. Right. The database can't mathematically know which one is the correct user 105. There's no way to automatically merge that. None. It usually requires manual, painful human intervention to untangle. And very often, data is just permanently lost. Which is brutal. So leader-follower is great until it isn't.

5:41 But there's another issue too, right? Geography. Ah, yes. The speed of light. Right. Because if I have a single leader in Virginia and it takes 200 milliseconds to acknowledge a right from a user in Sydney. That latency is going to absolutely kill your conversion rate. Exactly. So the obvious engineering impulse, especially if we want to avoid that single point of failure split brain risk, is to just put a leader in Sydney too. Right. Give the users local read and local write speeds. Which pushes us into multi-leader replication.

6:10 It does. And multi-leader solves your latency problem beautifully. A user in Tokyo talks to the Tokyo leader. A user in London talks to the London leader. The responses are lightning fast. So fast. The trade-off, of course, is that the leaders now have to asynchronously exchange their updates in the background. And immediately we're back to a synchronization nightmare. Completely. Imagine two users edit the exact same shared document or the same inventory counter at the exact same millisecond in different reasons.

6:38 You now have a hard conflict. So how do databases handle that? Well, the standard out-of-the-box conflict resolution strategy for most distributed databases is last-writer-wins or LWW. Okay. How does that work? It inspects the NTP timestamps of the two conflicting rights, kicks the one with the latest timestamp, and discards the other. Wait, hold on. Isn't last-writer-wins just a massive euphemism for permanently deleting user data? I'm... Because from an application perspective, both users clicked save. Both users got an HTTP 200 OK response.

7:12 They did. But one of them just gets their data silently nuked because of a fractional clock drift across cloud regions. I completely validate your horror here. Yeah. Because, yes, last-writer-wins is inherently lossy. You are explicitly prioritizing system availability over strict data preservation. That feels so dangerous for something like a financial ledger. Oh, you would never use it for a ledger. The alternative is writing custom application-level merge functions. Which sounds awful. It is. It means telling your code how to manually stitch together a shopping cart from Tokyo with a shopping cart from London.

7:45 It is incredibly tedious and highly prone to edge case bugs. There has to be a better mathematical solution than either silently dropping data or writing thousands of lines of bespoke merge logic. There is, actually. And it's where we turn to CRDTs. Conflict-free, replicated data types. Okay, CRDTs. I hear about these all the time. How do they actually save us from dropping data without forcing me to write custom merge scripts for every single table? The underlying mechanism of a CRDT is that it restricts your operations to mathematical properties that always converge, regardless of the order the operations arrive.

8:21 Give me a tangible example. Take a simple counter. Instead of a node sending a command that says, set the total value to 3. Which overwrites whatever state is already there. Exactly. Instead of that, the nodes simply exchange the operations themselves, like an array of increment-by-one events. Ah, I see. Because in math, addition is commutative. Exactly. It doesn't matter if you process A then B or B then A, the final sum is identical. Precisely. When the nodes sync, they just merge their arrays of operations.

8:49 Order simply doesn't matter. You can add CRDTs for counters, sets, and even really complex text sequences. Wait, is this how Figma works? Yes. That is the exact underlying mechanism that allows tools like Figma to handle real-time collaborative editing. Oh, wow. Yeah. Dozens of designers can drag elements around simultaneously, and the CRDTs ensure everyone's local state eventually converges to the exact same canvas. Without constantly overriding each other's work. Right. That's brilliant. Okay, so let's trace our logic here.

9:22 Leader follower gives us strict ordering, but introduces that split brain terror. Yep. Multi-leader gives us global speed, but introduces conflict nightmares. Unless you use CRDTs, yeah. Right. What if we just rip the band-aid off and fire the managers entirely? What if we build a system with absolutely no leaders at all? That architectural shift brings us to leaderless replication. The original Dynamo paper pioneered this, and you see it extensively in databases like Cassandra. Okay, so how does Cassandra do it?

9:51 In a leaderless model, the very concept of a primary node is completely erased. Every single node in your cluster is authorized to accept both reads and writes. But wait, if anyone can accept a write, how do we know a transaction actually committed? That's the trick. If there's no authoritative voice, no single source of truth saying, yes, this data is saved, how do you prevent total chaos? You enforce order through democracy. To make leaderless replication work, you introduce quorums. Quorums, like in a meeting.

10:22 Exactly. A quorum is a required majority of nodes that must acknowledge a read or a write for the operation to be considered successful by the application. Okay, so it's a board of directors. Go on. A quorum is like requiring a majority vote from the board to pass a resolution. Let's say you have a five-person board. Okay, five nodes. Right. If three out of those five members approve a new company policy, the decision passes. The write is successful. Yep. Now, if a week later you need to verify what that policy actually is, you don't need to track down all five members.

10:56 You just ask any three of them. Because of the overlap. Exactly. Because the original decision required three, and you are asking three, the math guarantees that at least one of the members you ask must have been part of the quorum that approved the original decision. That overlap is the secret. In database terms, as long as your write quorum plus your read quorum is strictly greater than your total number of replicas, You're safe. You are guaranteed that your read will encompass at least one node containing the most recent write.

11:25 And the beautiful thing for a back-end engineer is that Cassandra basically lets you tune this on the fly, right? It does. You can tune the quorum size on a per-query basis. Like demand a massive quorum for critical financial data, but use a tiny quorum for something trivial, like a user updating their profile picture. You absolutely can. But, and this is the big but, this highlights the hidden tacks of distributed systems. Latency. Latency. Whether you are waiting for a quorum of three nodes in Cassandra, or waiting for a specific number of followers to acknowledge a write and post-gress goal, you are fundamentally waiting.

11:59 Right. And waiting introduces the core choice between synchronous and asynchronous replication. And that choice essentially dictates your entire latency profile. It does. Synchronous replication means the node handling the write holds the client connection open, and it waits for the other replicas to confirm they received the data before returning a success message to the user. So your guarantees are ironclad. They are. But your write operation is now exactly as slow as your slowest network link.

12:27 Ouch. Yeah. If one replica is experiencing a localized traffic spike, your entire application's write request is just hanging, waiting for it. And the alternative is asynchronous replication. Right. The node stores the write locally, instantly tells the user success, and then quietly replicates the data to the other nodes in the background. It's incredibly snappy. The user experience is flawless. It is flawless. Right up until that primary node crashes a millisecond after telling the user success and before the background replication finishes.

12:59 Oh, wow. In that scenario, the write is permanently gone. Poof. Gone. Gone. The user thinks they successfully purchased the concert ticket. The UI showed them a confirmation, but the database system entirely forgot it ever happened. Okay. Because we are constantly juggling all these variables, latency, consistency flavors, inevitable network drops, we can't just guess our way through architecture. We need mental models to guide our tradeoffs. And the classic mental model is the CAP theorem. CAP.

13:28 Proposed by Eric Brewer. It states that during a network partition, a crisis where nodes physically cannot communicate with each other, a distributed system must make a brutal choice. You either choose consistency, which means refusing to answer requests because you can't guarantee you have the latest data, or you choose availability, which means answering requests even if you know you are serving stale, outdated data. You cannot have both. Exactly. But while CAP is foundational, it's a crisis model.

13:58 Right. I always think of CAP like having an evacuation plan for the building patching fire. That's a good way to put it. It's incredibly important to know where the exits are during an emergency. But network partitions in stable, modern cloud environments like AWS or GCP, they aren't everyday occurrences. Thankfully, no. CAP doesn't tell me how to run my architecture on a random Tuesday when the network is perfectly fine. That's why PACE-LC is the vastly superior framework for daily engineering.

14:25 PACE-LC optimizes for daily operations. It unpacks the acronym like this. During a partition, you must choose availability or consistency. The CAP part. Right. Else during normal operations, when there is no partition, you must choose latency or consistency. That's the real daily tradeoff. Even when the network is flawless, demanding strong consistency inherently costs latency, purely because of the physics of synchronous replication. Exactly. PACE-LC forces you to realize you live with this tradeoff every single day.

14:56 And you see this reflected perfectly in the technologies you use. Cassandra, for instance. Cassandra heavily favors the A and L in PACE-LC. Availability and low latency. Exactly. It uses last rider wins. It embraces eventual consistency. And it wants to give your user an answer immediately, even if that answer is slightly stale. Contrast that with Google Spanner. Spanner goes to absolute mind-bending extremes to choose consistency everywhere. Spanner is an engineering marvel. It achieves global linearizability by utilizing TrueTime.

15:30 Which is just crazy to think about. It is. TrueTime is an API backed by synchronized atomic clocks and literal GPS receivers physically bolted to the server racks in their data centers. They globally order every single transaction. They do. But they pay an explicit physical time penalty higher write latency to guarantee that certainty. And then you have a technology like DynamoDB, which is fascinating because it pushes the PACE-LC choice directly onto the engineer writing the query. Oh, yeah. DynamoDB is highly tunable.

15:58 By default, it gives you an eventually consistent read, which is blazingly fast and costs less read capacity. But if you add a simple flag to your SDK call, you can demand a strongly consistent read. It will take longer and it will hit the leader directly, but it guarantees fresh data. It forces the backend engineer to align the technology with the actual business requirements, feature by feature, rather than chasing abstract technical ideals. Which is exactly what senior engineering is all about.

16:26 Exactly. So to synthesize everything we've unpacked today, replication is just the mechanical process of copying bytes to survive hardware failure. Whether you design a leader follower, multi-leader, or leaderless architecture, and how you tune your synchronous wheat times or your core math, all of that simply determines where your system lands on the PACE-LC spectrum. Are you optimizing for speed or are you optimizing for absolute truth? Because the architecture is really just the physical manifestation of those choices.

16:55 It is. And that wraps up Chapter 2 of our system design deep dive. But we are far from done. Oh, we have so much more to cover. Because next time, in Chapter 3, we transition to the next massive bottleneck. What happens when reading from your perfectly tuned, eventually consistent database is still simply too slow? Because it happens. It always happens. What happens when the sheer volume of traffic demands a faster layer? Next time, we tackle caching patterns. We'll dive into cache aside, right through, and the dreaded thundering herd problem that routinely takes down major platforms.

17:31 That's going to be a fun one.