An Illustrated Proof of the CAP Theorem
mwhittaker.github.io
mwhittaker.github.io
E.g. often we have a fundamental tradeoff between latency and throughput, and it's impossible to get 100% in both metrics. However, we can still do very well in both, and that's what matters in practice.
We have the same thing with CAP. You can build a CP system with high availability.
I think the point is more about timing. If your client writes to G1, then you pretty much have to wait until the write has propagated to the entire network to acknowledge the write or accept some risk that some other client will read back stale data after the write has been acknowledged. I should think this would be obvious to coders also, but it is not at all obvious to managers.
I say "pretty much" above because there is technically a ridiculous third option where you essentially lock down the whole network with every read. But it still doesn't get you around CAP.
The first part, PAC, is your traditional CAP theorem - in the presence of partitions (P), you can provide either availability (A), or consistency (C). The second part, ELC, describes the system characteristics during the normal, non-partition case. It reads as "else (E), you can provide either low latency (L) or consistency (C)".
Even tho apart from some curious outliers most systems are either PA/EL or PC/EC, I find the framing helpful for reasoning about a system in more than just the partition or failure case.
"If your client writes to G1, then you pretty much have to wait until the write has propagated to the entire network to acknowledge the write"
PAXOS, for example, just needs it to be written to a majority of nodes.
If you write it to node "key % n", where n is the number of nodes, then only that single node needs to be up & available to all writers / readers.
Hmmmmmm.
Of course you can prove something exists by an example. And you can prove a property does not always hold with a counter example.
But if you want to prove that something doesn't exist, or a property always holds, then you shouldn't be looking for examples.
> Proof by example or analogy is not a proof.
So as you also said "Of course you can prove something exists by an example".
To answer this: > Proove that you can't walk. Don't walk.
As I said, proof by example works only
> to prove that something exists or that a property is not valid for every element of a set.
Divide your lifetime in seconds. With that, you're trying to proove with the property "walk" isn't valid for every element of the set lifetime. To proove that the property can't walk dosen't hold for every element of the set, just show an example of second while you were walking and you're done.
[1] https://en.wiktionary.org/wiki/Reconstruction:Proto-Indo-Eur...
... which is akin to Latin fugio, fugere -- to flee. The analysis of OE begietan as be+get strikes me as faux-etymology fashioned after behold and the like, where beholden has little to do with holding, all the same.
This is a hobby of mine to the point of becoming neurotic. It's all in vein if a play of words needs explaining. It's still kind of insightful, if you'll entertain me a little longer.
If bʰegʷ- means to run, flee, then how is bʰeg-, to break, related? Why is to break rather reconstructed as bʰreg-. Why is bʰegʷ- given an alternative form bʰewg-? Does preḱ-, to ask, fit in here; or prey- or what it was with a related meaning? Maybe wreg-, whence wreck?
What about bag, pray, pay, fag, vag, way, weigh, etc etc.
Maybe to break away from always meant ''departure'', ''to go, run away''. Why are hurried news breaking? Why does German have "brandaktuell" instead, burning news? I have to break it to you, I don't know. I can already hear you begging the question to stop. But I have one more. Is a beggar someone living out of bags? I really don't want to know.
A proof establishes truth, it does not typically explain. And yes, some assumption was wrong. The assumed premise was wrong. A system which is consistent, available, and partition-resistant does not exist. The truth of the theorem that no such system can exist is proven. I imagine it's not the proof you dislike, but the article about it which you wish had more explanation which is reasonable. It would be a different matter to discuss which solution is "best possible", though, as that would vary greatly on the application. If you're running a bank, you would never want to sacrifice consistency. If you're running a blog, availability and partition-resistance is probably more important.
Banks do sacrifice CAP consistency. As it doesn't mean you don't have consistency at all or anything like that. Just that you don't wait for writes to become visible to all nodes.
Then in the example, client writes to V1, V1 waits T time to communicate with V2, it doesn't manage to so it returns failure to the client and kills itself. Then the client writes + reads to V2 instead.
Does this break availability because returning an error that you can't act on a client request counts as ignoring it in this informal definition?
Your example can't deal with there being no path at all between V1 and V2. So even the client can't reach V2. In that case the entire system V is unavailable to the client.
Availability in CAP is about being able to reach a node, but not all of them and because of that the system doesn't function. I guess killing a node when a problem occurs can be argued to defeat CAP (all the nodes that are up can reach each other), but it definitely doesn't improve the situation.
Also what does you system do when V5 out of a 5 node cluster crashes? Do the nodes V1-V4 kill themselves, because they can no longer reach a higher ID node?
TBH my 2-node system was a simplified version of Corosync which I've worked with a bit, in that you allow quorum if you have half+1 of nodes up - so 4/5 up will maintain quorum. The highest-node lives rule is just a tie breaker in the case of a half/half split which the 2 node system always has when one node goes down. Good point though, that rule alone is definitely not a good idea.
In practice, partition tolerance is the most important property to have, since there's no point in having a distributed system if it can't function if any part of the system is down (it has to be fault tolerant).
>usually because of networking problems.
In distributed systems, those are the same thing.
If I am a node -- n1, communicating with various other nodes, the only way I can tell if another node -- n2 is alive or not, is sending messages to it and receiving responses. If the network is having issues and I can't communicate with it, for all intents and purposes, that node n2 is down from my perspective n1.
Edit: Also, diverging of state relates to consistency and not partition tolerance.
You can get partition tolerance at the expense of availability,for by refusing writes during a partition.
This is really a question of definitions, no?
If partition-tolerance is a defining property of a distributed database system, then of course they all have it. Plenty of database systems don't have partition-tolerance: non-distributed database systems.
Non-distributed database systems are inherently partition-tolerant because they can not break into partitions in the first place.
You only really want to sacrifice partition tolerance if the precise correctness of your data is not tremendously important. A CA system can always respond and appear correct, but network partitions result in desyncs. I believe this approach is fairly popular in video games, where a CA system for multiplayer can provide each client with a coherent gameplay experience at the expense of things just breaking if a partition forms. Sorry, game over, please try again.
Most serious systems, though, require P.
One potentially interesting example that uses the current fashion of the day could be to compare payment networks. Visa's distributed payments system is CP, which is to say that it's a consistent system and getting a successful response back means your payment has gone through, but you have no guarantee of getting a successful response. The bitcoin network is AP, which is to say that it's a highly _available_ system which will never fail a response (so long as you can reach one node), but getting a successful response is no guarantee that your payment has gone through.
In all cases, the sacrificed letter is not gone, merely imperfect. This does not mean that Visa is not Available, it just isn't _perfectly_ Available. It does not mean Bitcoin is not Consistent, it just isn't _perfectly_ Consistent. It does not mean that your video game is not Partition Tolerant, it just is not _perfectly_ Partition Tolerant.
You can sacrifice any of the three and still have a distributed system, it will just behave differently.
There's no point in having a system if it gives incorrect answers (inconsistent).
In the real world, you have to make tradeoffs.
Partition tolerant is "most important", it's awkwardly factored out of Availability -- "tolerant" is a variant of available. Your system can't be "partition intolerant "just because that means "not always unavailable".
Another way of looking at it is what you imply: "partition tolerant" is a synonym for "distributed".
(As in the formal definition is not what we actually need to get a consistent and reliable system for our real-world application. We need something a little less and then suddenly it's possible to get a similar solution that actually provides C A and P.)
However his definition reads: "every request received by a non-failing node in the system must result in a response"
It doesn't say anything about the timing of response. The network partition will be resolved eventually and thus the write operation will complete. Or it can time out and return an error, which is also fine based on the definition of availability (must result in a response).
The way I saw the proof go is: to satisfy CAP, you need replication, and you need high availability. To be consistent you need replication. But if you don’t have a reliable network, you don’t get reliable replication. Therefore you can’t be consistent. To me, this is a trivial result, and I’m not sure that this is what the CAP theorem intends to say?
It would be better to do, as you say, assume that a system is CA, And prove why it can’t be partition tolerant. Saying “network is unreliable” isn’t a suitable jump in argument IMO.
I think that might also make split brain situations more likely if all the client brokers didn’t have the entire cluster topology known. You’d end up with different data in different places if the clients are partitioned in addition to the brokers.
disclaimer: i am not well educated on the literature in Distributed systems :(
There are various ways to reduce the impact and probability of partitions, but increasing the quantity of nodes does not make partitions impossible -- in fact, you now have to worry about n possible partitions (n being the number of nodes).
the only way to guarantee partition tolerance is by majority votum. If you do that, you can no longer guarantee availability. For example, a cluster of 11 nodes might be partitioned thrice. (2 + 4 + 5) None of the nodes are allowed to answer, breaking the availability guarantee.
So instead we talk about high availability. All sorts of things cause downtime besides network partitions. If the network is reliable enough, then you can still achieve high availability.
>The purist answer is “no” because partitions can happen and in fact have happened at Google, and during (some) partitions, Spanner chooses C and forfeits A. It is technically a CP system. We explore the impact of partitions below.
...
>Conclusion
>---------
>Spanner reasonably claims to be an “effectively CA” system despite operating over a wide area, as it is always consistent and achieves greater than 5 9s availability.
From the paper
Lots of DB vendors have tried to circumvent the therorem by introducing their own concepts or stretching the definitions. "We can beat the CAP theorem"
Usually companies find this out the hard way, with catastrophic loss of business data.
The "CA in practice" claims are purely marketing.
[1]: https://storage.googleapis.com/pub-tools-public-publication-...
There are systems tolerant to these partitions (see PAXOS)