The CAP theorem. The Bad, the Bad, & the Ugly
blog.dtornow.com
blog.dtornow.com
Newtonian physics/mechanics: good enough in a lot of cases.
Einstein accurate, but unnecessary in most cases.
In many cases CAP is good enough for us to have the conversation about how the system works. One can then formulate plan for when it doesn't. The fact that it's imperfect at a formal level is academically interesting, but technically irrelevant for a LOT of conversations where it has utility.
> In theoretical computer science, the PACELC theorem is an extension to the CAP theorem. It states that in case of network partitioning (P) in a distributed computer system, one has to choose between availability (A) and consistency (C) (as per the CAP theorem), but else (E), even when the system is running normally in the absence of partitions, one has to choose between latency (L) and loss of consistency (C).
People who want CAP to be profound and useful are missing the point - it just tells you some things that you might want to try to achieve are provably impossible. Just like Newton’s first law tells you things can’t accelerate without a force (duh) and the first law of thermodynamics tells you you can’t get energy out of a perpetual motion machine (duh) and the pigeonhole principle tells you if you put n things in less than n groups, at least one group has more than one thing in it (duh); CAP tells you you can’t build a distributed, partitionable data system that is 100% consistent and 100% available (duh).
It’s not that deep or profound but it is proven and it eliminates a whole class of ideas from meriting further thought because they’re demonstrably impossible. That lets us get on with making the most of what is possible.
(by the way I think you have your Newton's law order wrong, that's the second one you're referering to)
Second law gets more specific, and clarifies that the size of the acceleration is proportional to the size of the force.
There will occasionally be network partitions. When there are, a given node can either respond (potentially inconsistently) or not. So you pick some balance between consistency and availability. Latency is really just a proxy for availability - as latency tends towards infinite, availability tends towards zero. Of course you can wait until the network partition is resolved and the nodes are caught up - but I think it's simpler to consider that as not being available for a period of time rather than as having a high latency at that time.
Consistency or Availability - you don't have to pick one, but the more consistent you want to be (in the case of network partitions) the less available you'll be.
Failing to understand the implications of one over the other is often sign of immature architecture.
There are databases that can provide high consistency and availability.
Sure, as soon as you decided to distribute your system across a network you opted into a world where partition can happen, and you will have to give up consistency or availability.
Mainframes, though, provide consistency and availability by being unpartitionable except through use of a chainsaw.
A nuclear blast in northern Virginia can take out all the availability zones for AWS US-East-1 - that doesn’t invalidate the availability claims of a system that is distributed across three AZs, it just puts an upper limit on it.
Though, thinking about it more, what it also tells us is that we need to be clear about what "the system" is. There are, after all, some systems that _do_ take nuclear blasts into account, or systems that take "physical breach by a hostile actor" into account.
Less dramatically, we make trade-offs all the time as to what we consider "in" system vs "out", and it's good to be conscious and explicit about them. We see this a lot in UI, especially web UI in terms of things like what browsers to support, whether we care about users' battery life or data cap limits, if we foresee ourselves running the system in a jurisdiction with different regulations around data collection and usage, etc.
If so, I agree - in that case it would be similar to running two severs in the same AZ. Which provides a certain form of availability.
Otherwise I disagree, because as soon as you have a single point of failure, your system is limited by the availability of that specific part, whereas in a distributed system you can by design increase availability by adding more nodes, at least up to a certain point.
Think about it: I can run the two servers I mentioned about in the same AZ. But I can also move them to a separate AZ each. For a distributed system it doesn't matter conceptually. That doesn't mean that there is no impact on the availability or performance, it just means that I don't have to change the logic of my system.
For a mainframe however that's not true - and that is what makes the charm of it: because I don't have (and don't want) to care, which simplifies things a lot; at the expensive of (the option of) the availability of a distributed system.
> One, quite commonly known, is that two separate mainframes can run completely in lock-step if they are less than ≈ 50 km / 30 mi apart
Since they are 50km apart, they are connected to each other by one or multiple cables. I assume that they are also both connected to the internet (or something else that matters) so that if one burns down the other can take over. Correct so far?
If so, then please answer this questions: how does the system behave if all the cables between the two machines are severed, for an unknown time?
The ultra-scalable techniques that Google, Facebook, ect, use are great, but most applications do not need that kind of scalability. Throwing a little more money on a fancy mainframe is much, much cheaper than the programmer time needed to build a distributed system.
Which is odd, when you think about it, since you allegedly made this choice in order to get better availability.
Also, mainframes are pretty big, multi-processor machines. They have their problems, but having all those components give out all at once? I'm sure it happens, but I've never heard of it.
Why that? I can distribute my nodes over multiple datacenters. Netflix does that even with whole regions (to my knowledge) to avoid downtime when a region goes down.
But yeah, the theory is what I'm talking about. Not putting the theory into practice is a whole different matter. :-)
For example, if you need to account for cosmic rays, you can't prove anything meaningful about software. Redundancy won't get you out of this: you add redundancy, I'll add more cosmic rays.
All that does is force a useful analysis into a muddy and probabilistic one. It's actually pretty important to account for cosmic rays! But you do that on top of an analysis which presumes that the hardware performs correctly. It's not a useful reality to expose to a proof assistant.
To anticipate an objection: yes, you absolutely can rule out network partitions in an analysis, if that's a useful thing to do. For a packet-switched network, it isn't useful. But different network topologies exist, or at least used to, ones where once a circuit is negotiated, you can say useful things about the network without accounting for the kind of equipment failure which might break that circuit. For packet switching, you're going to have a really bad time if you don't account for partitions, because those are expected behavior, one of the properties of the system under consideration.
> as soon as you decided to distribute your system across a network you opted into a world where partition can happen, and you will have to give up consistency or availability.
So we are talking in the context of distributed networks. Then:
> Mainframes, though, provide consistency and availability by being unpartitionable except through use of a chainsaw.
using the words "consistency" and "availability" here is clearly refering to the same words in the sentence before. Hence we are still in the context of distributed networks. And therefore you cannot just rule this out as an "act of god".
Had OP said "A mainframe has high availibility compared to my laptop and for me that is more than enough." then I wouldn't have said anything.
Not only that, a weird network partition is highly rare and I have never seen it. In most of the cases, there is only one network and either the server is up and connected to "the" network or is not connected to the network. And I believe engineers have tendency to overengineer for this downtime and not thinking about much more probable ones.
I have asked many who says we need three nodes and three pods for each service minimum, and no one could answer why.
1) to ensure you can deliberately take one offline for maintenance and still have redundancy in case a single node goes offline
2) in some systems an odd number of nodes is needed to ensure no ties in leader elections or decision votes. Three is the smallest odd number that has any redundancy.
2) 1 is odd and could win the election;)
Running three nodes is a rule of thumb, not a hard rule for minimally guaranteeing availability.
No, they can't be rebooted any time. Where are you getting this information from?
Note in particular:
> For instances that launched from an Amazon EC2 Auto Scaling group, the instance termination and replacement occur immediately
Interesting. I had seen EC2 instances running for multiple years and never saw this so I assumed that this isn't possible. I know for a fact that GCP has live migration where while the service would be degraded for few seconds, you don't need to do anything so I assumed AWS also had something similar.
Some maintenance person disconnected power to the two servers, then reconnected the power assuming they would reliably recover. The admin instructions were to never do this - always boot up in sequence.
They booted up simultaneously, and the fully redundant network switches took too long to boot fully (Cisco), but started passing internet traffic. After a timeout the servers each assumed they were the master in a degraded cluster, so proceeded to make divergent modifications on the DRBD replicated storage and to serve requests.
I never found out why this happened in spite of the direct ethernet links between the servers which they were supposed to use for synchronisation decisions, but it did.
Recovery required manually comparing changes in files and databases to decide whether to merge or discard.
This problem was avoidable but an adequate fix evidently wasn't in place (mea culpa, limited time and budget).
It did not help that Pacemaker+Corosync was used, before Kubernetes was popular, and Ubuntu Server shipped a very buggy alpha version of Corosync that corrupted itself and crashed often, despite upstream warning it was an unreliable version. I had to manually build a different version of those tools from Red Hat source, because it was too late to change distro. This is one of two reasons I don't recommend Ubuntu Server in professional deployments any more, even though I still use it for my own projects.
Three servers, or two servers and a third special something for arbitration, is a standard solution to this problem.
But it's only useful for a stateful distributed system, like a database or filesystem with some level of multi-master or automatic failover.
There's no need for three nodes or any particular number, for stateless nodes like a web service whose shared state is all calls to a database or filesystem on other nodes.
Technically you don't need three servers. It's enough to have a cheap component or low-cost tiny computer to arbitrate. Even sending commands to the network switches to disable ports (if the switch doesn't behave too strangely, as the Cisco switches did in the above!), or IPMI commands to the other server's BMC. Just about anything can be used, even a high latency, offsite tiny VM, as it isn't needed when the main servers are synchronised.
Notwithstanding, I still find it very useful as a general guide when designing distributed systems.
We generally accept that clients may have to perform retries, definitely jittered and hopefully bounded. Given that's the case, what does it matter whether a single instance takes 5 seconds to come back online after failure (possibly rescheduled on a new node), or multiple instances take the same 5 seconds to recognize leader failure and elect a new leader? Sure, they're unlikely to be identical numbers, but say they're within an order of magnitude and both have error margins so it's a wash.
An architecture astronaut will put forward a design with complex leader election, endpoint discovery, distributed locking, strong consistency, etc. for the latter solution and pat themselves on the back for their expertise and professionalism. They'll waste a lot of time for both dev and SRE as long as the service lives.
A reasonable person should be able to acknowledge that, given either solution will have the client retrying about the same amount of time, the simplest solution that delivers that experience is sufficient from the client's point of view and vastly preferable from a maintainer's or operator's point of view.
Of course not every scenario will be like this. Sometimes startup is unavoidably a lot slower than re-election. It's just rare you see this evaluated in terms of client-facing numbers before a very costly architecture decision is made.
I find that reasoning about consistency is much easier from that perspective, for example it immediately gives you an intuition on why individual CRDTs work.
Forgetting for example is fine, but deletion is not.
No, it's two out of three. CA systems don't use a network. Consider a database and a work queue which communicates with the database, where they run on the same server. You can achieve consistency and availability in such a system, but only by eliminating partitioning, and there's only the one way to do that: don't run the systems on a network. Nor is this an unrealistic architecture! Far from it, it is in fact the one you should choose if you need both consistency and availability.
The parts of the article about the CAP theorem and its consequences are on fairly solid ground as I see it. But the observation about the CAP conjecture is bootless, the conjecture that CAP is "pick two" holds up to scrutiny. Not that it proves the conjecture, just that all three of the "pick two" options describe meaningful systems.
That is people doing implicit or explicit marketing of some enterprise data landgrab telling you things like they can scale arbitrary SQL joins across distributed tables. Because large data sets do not teleport across the wire for merging.
Foundation db, Cassandra, and kinda dynamo (I don't trust their new global replication) actually scale, but I've never used gcps db techs.
CAP is imo one of the best distributed systems principles ever stated. It is simply state, compare to say the nitty gritty of say paxos vs raft. It is a very useful first principles exercise for deconstructing whatever harebrained claim some marketroid is foisting upon you.
Is that why S3 works so well?
1. Get write request to a node
2. Node sends lock on that updated data to rest of nodes
3. After they send back OK the write happens
4. Propagate to rest of nodes
5. After receiving OK on the update from all of the nodes send notification to nodes to lift the lock
So if a partition happens the system fails if it is CA while with e.g. CA it cannot guarantee strong consistency without resolving the split brain (killing of one of the partitions based on e.g. quorum)?
If you need stronger primitives than atomic read/write (e.g. CAS), then ABD is insufficient and you need consensus in the asynchronous model (although synchronized clocks can allow even CAS to be implemented with ABD + leader election with leases[2]).
[0] https://dl.acm.org/doi/pdf/10.1145/200836.200869
[1] https://www.cl.cam.ac.uk/teaching/2223/ConcDisSys/dist-sys-n...
Invariant confluence, determines whether an application requires coordination for correct execution.
You can tailor the invariants to your requirements making invariant confluence a much better tool than CAP
- CAP came up
- all parties immediately agreed that CA was impossible
I still think it’s a useful model for forcing people to think about their failure modes though. Systems, largely, tend to fall into either CP or AP or be horrendously expensive.
You have to choose your trade offs.
I wish the author had expounded a little further on better options. I’ll read the paper linked, but there’d be more punch with more inline content.
[[The “Pick 2 out of 3” interpretation implies that network partitions are optional, something you can opt-in or opt-out of.]]
is you can choose to not have a distributed system.
The approach we took was to have a single database node be the primary and another replicating from it asynchronously. When a node went offline, a third monitoring node made a decision to promote the backup to primary and communicate this decision to all clients. The three nodes were located in three different cities and as such the possibility of all three having issues at once was reduced. By doing this, we could literally pull the plug on the primary node and within 30 seconds the clients were talking to the backup-now-primary. If the monitoring node failed, of course we had no way of automatically switching database nodes but again it had to be a major issue for a node in Seattle to get cut off at the same time as the nodes in Dallas or DC. This worked because of the kinds of data we processed so that our 30 second failover time was acceptable AND a small window of lost writes was also OK.
I think another interesting approach is to have your nodes explicitly communicate write status when responding. Basically when you write(node1, key, val) and expect node1 to normally propagate the write to node2, in the case of a network partition node1 would respond with ack(key, nodes_written=[node1]) explicitly excluding node2 from nodes_written if it was unable to push the write to node2 in a timely fashion. Similarly, read(node1, key) could have node1 return resp(key, val, inconsistent=true) or again spelling out which nodes do or do not know this value.
This would give the application a way to decide if the write should be considered successful or not based on what key represents. For example, updating the current position of a fast moving object that sends updates all the time could easily lose a write or two without it being a problem for the user, but a financial transaction could not.
Lastly, the approach I never explored but was curious about is the idea of the client being responsible for pushing all values to all nodes rather than replication happening in the background. This would effectively allow the write to happen simultaneously and also know if the write was a fail, a partial success, or a full success. Then a read could be done from just one node, but before returning it would poll other nodes and return the value that is consistent amongst the majority of the nodes or a failure if it could not reach any of them.
You can client replication at industrial scale by having clients push writes to Kafka and then have multiple, independent systems read and apply them. You still have a SPOF on Kafka of course. Another practical issue is that if you take this approach you'll likely need a utility to detect and correct data drift. That's been a feature of the systems I've seen that use this approach successfully.
It's still extremely useful as a mental model for tradeoffs in designing a system.
All it's telling you is there's a pendulum from availability to consistency and generally. The more you want of one, generally the less you get of the other, so choose the amounts you want of each deliberately.
People take things too literally and seriously. Has anyone in their company really argued that, "we should choose availability instead of consistency"? It's way more nuanced than that, and as far as I know, everyone knows that
… I'm usually experiencing the CP/AP split from the seat of a user, using the system. From where I sit … yes, it certainly feels like people are, somewhere, saying "we should choose availability", given the number of systems I've had to interact with that are trivially not CP.
Specifically, refer to https://jepsen.io/consistency — these are better, more specific terms, IMO; the number of systems that don't obey "Read Your Writes" (at the very bottom of the tree!) is pretty stark. Many Azure services, for example, are not read-your-writes; this means I'm perpetually wrapping things in loops that attempt to wait for the upstream system to become consistent, which is impossible to do in any manner that's foolproof, prior to moving on to the next API call (which would otherwise fail, if it depends on a write from the prior API call, but can't read it, because read-your-writes). Off the top of my head, I've seen empirical violations of read-your-writes in all of ARM, AAD, ACR. My latest container registry is also not read-your-writes.
S3 used to fail read-your-writes, in certain circumstances. (That have since been fixed.)