Please stop calling databases CP or AP (2015)
martin.kleppmann.com
martin.kleppmann.com
I always found it frustrating that the semantics of the word "Consistent" in ACID has nothing to do with the semantics of CAP "Consistency". ACID Consistency refers to integrity constraints while CAP Consistency seems to refer to Cache Coherence across nodes in a cluster. I was a little surprised that Kleppmann uses the idea of "Linearizability" to explain the discrepancy but it does clearly identify the same mistaken assumptions.
I was also frustrated with the NoSQL movement's focus on Partitioned clusters when my intuition was that it was a very rare failure mode. Perhaps we can now return to the metrics of data loss (RPO - Recovery Point Objective) and Recovery Time Objective (RTO) since we seem no closer to fully transparent failure recovery.
Regardless, the goal was always to have available and scalable databases that operate on clusters of commodity cloud servers and all of the SQL, NoSQL, and NewSQL solutions have converged around that goal. The transition to cloud computing seems to have been more about NoRAID than NoSQL. Some technologies like LSM-Trees address important use-cases in commodity clusters so maybe we/I shouldn't focus on the truthiness of the original rhetoric.
Network errors are very common among errors database has to deal with, and particularly common when we are talking about networks between datacenters.
Another part of it is that all hell breaking loose can sometimes look very similar to clean partitions, depending on the network. When network QoS classes are available, what appears to be increased latency or packet loss for one class may be 100% packet loss for the other. Though, you're totally right that all hell breaking loose can be very different - the set of systems that handle flapping & other ugly connection failures well is smaller than the set of systems that handle clean disconnections well :)
For instance, a metro might be isolated from the rest of the network due to simultaneous fiber cuts.
https://www.bbc.com/news/technology-50851420
Or we could see heavy, prolonged network congestion on the backbone, impacting comms to/between backend services, while frontend serving is unaffected.
https://status.cloud.google.com/incident/cloud-networking/19...
(Incidentally, IIRC, Youtube's video serving infra was unaffected, but the backend jobs for youtube.com were largely inaccessible.)
Or improperly-executed maintenance that isolates a metro by draining routers in the wrong order. Insufficiently-provisioned network capacity (or network augments lagging far behind). Backbone traffic engineering gone awry. Underwater fiber cut due to shark attacks. QoS inversion (i.e. user-facing traffic as high priority, but replication as low priority) causing a partition during periods of heavy network load.
While yes, a lot of these scenarios might have (and did!) impact Google's serving capability, Google servers were still IP-reachable. 8.8.8.8 should have still been able to serve DNS (albeit with potentially stale data).
(In more recent news, if you use a transit provider to provide connectivity between your datacenters, and said transit provider suffers a multi-hour outage, that's a clean partition.)
You're more likely to see it when you are running services across multiple datacenters/availability zones/regions/whatever. It still happens frequently enough that it's a scenario worth keeping in mind. You probably have your application servers local to your database servers in that other datacenter, so if your public facing network falls over in a datacenter, there's probably still some amount of work happening that modifies things in the database - what happens when that connectivity is restored?
Being able to gracefully handle network partitions is hugely important in distributed computing.
PACELC is an improvement over CAP in some sense because it is a finer-grained characterization of distributed systems, but it still suffers from the problem that the stated properties are too strong. As Kleppmann explains in the article, there are systems that strictly only fulfill the partition tolerance part of the equation, yet are still very practically useful. This is because while they may not satisfy strict linearizability and availability, they have softer guarantees (e.g. defining "availability" as "99% of requests complete in 1 second") that we can rely on in practice. This reasoning extends to PACELC as well.
I am inclined to agree with Kleppmann. CAP, and by consequence PACELC are useful mental models in the extreme, but their characterizations are really strong. It seems to me that e.g. probabilistic guarantees based on actual failure statistics are much more useful for the real world, because they aren't as strong as CAP/PACELC guarantees.
Edit: I am also not sure if I agree with calling CAP "dumb". It's a very useful theoretical result in a particular model. I think the problem is when people try to extend the implications of the theorem beyond the restricted model it operates in.
In a different arena, I also know of no better book at the interface of distributed systems theory and practice than Kleppmann's. It's useful for engineers while being enriched by deep references to the literature. If that doesn't demonstrate mastery, I don't know what does.
The book chapters do an solid job laying the ground-work for those papers. The depth is in those references. Read them if you can!
The article claims that the CAP theorem only prohibits linearisability as a consistency model. The original CAP paper talks about "Atomic consistency" which is clearly at least linearisability. However as the proof formalizes, it's obvious that if you accept writes to totally split and separated systems which cannot communicate together (network partition), you cannot be consistent across the two parts. This is confusing because this kind of consistency only loosely matches the definitions in the ACID acronym which are geared around invalid database states in a single machine database, not about different requests to different database machines giving inconsistent answers.
Linearizability versus serialisability (and strict serialisability) are explained well by Peter Bailis http://www.bailis.org/blog/linearizability-versus-serializab...
I did not realize that consistency and atomicity also are overloaded between databases and distributed systems. That certainly adds to the confusion! I must have missed it when reading the original article.
More real-world experience is also really valuable, and an application-first approach to handling distributed systems works really well. I guess my hope is that eventually the theory becomes robust, matching the real world enough so that we won't have the problems of mis-applied theorems. IMO the best theories are intuitive and help us expose holes we could not catch in the real world.
If you can warranty high availability (by other means) then you can focus on consistency.
It also says that you should design your data so that becoming partitioned is not a problem.
What it doesn’t do, particularly, is help us build them.
Why?
> Despite being a global distributed system, Spanner claims to be consistent and highly available, which implies there are no partitions and thus many are skeptical. Does this mean that Spanner is a CA system as defined by CAP? The short answer is “no” technically, but “yes” in effect and its users can and do assume CA.
The point is that if my application will actually go down anyway if the WAN craps out, then really, I don't actually need P. The application would also be significantly simpler if it assumes that if the application can work, then the DB is also up. And one seems that realistic systems built on Spanner simply assume CA and then get on with life...
Can you guarantee no further modification to the data in the database occurs when the WAN craps out? What happens to in flight requests, batch/scheduled jobs, etc.? There's no situation where your application continues modifying the contents of the database even when there's no WAN link?
I've seen very few real world services where there is a 0% chance of data changing just because the application servers can no longer talk with clients but can still talk with the database. What happens when those databases try to rejoin the cluster and their data is now inconsistent, if you did not take into account partition tolerance?
Often you can just let the people choose whether to accept the latest state of a particular partition of the database, or wait until it becomes consistent again.
CAP theorem is older than the career of nearly everyone on HN.
This is not a snarky comment, I really do not understand it. (But extremely off topic, I admit)
This isn't funny?
These days, I think we expect everything to go to extremes. It has to be hilarious, not just funny. It has to be tragic, not just dramatic. It has to be action packed, not just action.
Seinfeld wasn't hilarious. It was just funny, in kind of the same way your friends are funny when you're joking around the table during poker, but any outsiders would probably think you weren't funny at all.
And if someone don't understand it then it usually does not end well if they work with distributed data stores.
I know this is a useless construct in practice, and there is probably a flaw in my reasoning, but it seems to me that you have to establish a non-infinite timeout for the proof to be consistent.
CAP tells us that databases are fundamentally broken, Godel theorem tells us that math is broken. But besides being important theoretical points, they don't matter that much in practice. In the same way, the two generals problem doesn't prevent TCP from working fine.
Gödel's incompleteness (that, by a diagonalization argument, formal systems of minimal expressivity will contain true statements that cannot be proven) applies and will continue to apply to all such formal systems. The bar is so low that any interesting formal system is affected by it. Also, it does not make math broken at all - those formal systems can and in fact must be consistent for the proof to go through (1st incompleteness theorem), and it's just a pity that whatever formal system one comes up with, if it contains arithmetic, it won't be able to prove it's own consistency (2nd incompleteness theorem).
In the words of von Neumann: Gödel's two past papers ("Widerspruchfreiheit" and Continuum Hypothesis), are in any case worth more than the total literary output, past, present and future, of most mathematicians [...]