CAP twelve years later: How the "rules" have changed (2012)
infoq.com
infoq.com
The “P” in CAP means “tolerance of partition events”. But that’s really not what most engineering teams are thinking about when they talk about CAP theorem. Teams that really want tolerance to network partitions are really talking about building parallel instances, for example, cloud providers implementing multiple regions. If they need some level of consistency across those instances, then really, it becomes a discussion about latency.
If you think about it, what you’re really talking about with the network partition events is how stale or latent the data on each side will become, and whether or not, it will become consistent again, particularly without human intervention.
When you talk about it in terms of CAL, now all the properties are continuous, and you can discuss the trade-offs to each property with respect to network partitions.
Klepmann makes an interesting critique [3], that neither is particularly useful because both imply that you must either pick linearizability or total availability. When in practice, there have been plenty of successful databases that present neither characteristic, and thus are neither CP nor AP.
[1]: https://www.cs.umd.edu/~abadi/papers/abadi-pacelc.pdf
[2]: https://dbmsmusings.blogspot.com/2010/04/problems-with-cap-a...
[3]: https://martin.kleppmann.com/2015/05/11/please-stop-calling-...
When too many nodes fail or are unreachable, do you sacrifice consistency, or availability?
And even then it's misleading, because you can still be consistent and partially sacrifice availability (allow local or mergable updates) or be available but partially consistent (allow updates, but have potential conflicts to resolve, maybe manually, when the partition is healed.)
You can even make different A <-> C tradeoffs per request/transaction.
Distributed system are complex, who knew?
It takes time for all nodes to come to agreement and in that time I am able to request _technically_ out of date info from a node.
Consistency is often the one CAP that isn’t prioritized because if the network isn’t consistent for a few seconds, the application probably still works and that’s the case for bitcoin.
The network is “eventually consistent”
Don’t forget the nasty bit about partitions: if I think nodes 1-5 are reachable and the rest are unreachable, I can’t assume that the rest of the nodes are down or idle — some other client may think that nodes 1-5 are unreachable but nodes 6-15 are reachable.
Layer-3 IP Clos networks are pretty much a standard within a datacenter, making multiple redundant paths available. Across datacenters, optical circuit switched networks with lots of redundant paths are very common.
Modern network fabrics also have end-to-end mechanisms for bandwidth engineering based on per message QoS/priority set by applications as well as very accurate time-keeping.
Observability mechanisms and operational practices have also tremendously improved. A single node becoming unreachable from all other nodes is quite possible. But a single node being reachable to some subset of other nodes in the cluster and not others is nearly impossible. Only way a network partition can happen is due to misconfiguration. Purely layer-3 IP clos networks with BGP advertised multiple ECMP routes are a lot simpler than layer-2 networks of the past. With modern switch operating systems, modern devops practices, Infrastructure-as-Code etc. misconfiguration is nearly impossible.
At a higher layer, disaggregated storage and compute architectures are standard for any distributed database system. In this model, storage fabric is quite simple and much more distributed and redundant. Today, it is a lot easier to guarantee a consistent and highly available transactional write to a quorum of simple storage nodes spread across multiple datacenters in a region. On top of that, building higher-level isolation guaranteeing database transactions in a tier of ephemeral compute nodes is relatively much simpler.
This kind of system is extremely resilient to failures at all levels – node failures are detected and lost processing capacity recovered in seconds; failed transactions are retried in milliseconds; loss of redundancy due to failed storage nodes are rebuilt in seconds; loss of network capacity due to failed circuits or intermediate switches are routed around in microseconds.
So, in today's context, we build databases that provide flexible transaction isolation guarantees that is selected appropriately on a per transaction basis by the application. Application also specifies the latency budget it has available per transaction. With all these improvements, thankfully, application developers don't have worry about CAP theorem like they had to in the heydays of NoSQL databases in 2012.
Due to the quorum requirements, probability of failure increases with larger quorum sizes, as the number of intersections (and therefore potential points of failure) increases with more complex quorum arrangements.
Well, that's not true due to two key design points:
- Size of the read or write quorum can be chosen on a per transaction basis depending on what level of availability and consistency you want. The design goal is to lower the likelyhood of the minimum number of quorum members you need for your transaction will not be available.
- Given the transaction processing nodes are ephemeral and independent from storage nodes, if a transaction processing node which is designated as a member of a quorum dies, it can be quickly and instantly replaced with a new one and all in-flight operations will be retried by the smart clients.
Overheads and latency issues in node replacement, particularly under conditions of unexpected high load or network difficulties, compound these challenges. These issues often manifest in correlated failures, which are more common than isolated node failures.
In this landscape we are still in compromise and trade-offs territory. I would refer to these two papers as insightful demonstrations of these challenges:
"Consistency vs. Availability in Distributed Real-Time Systems" - https://ar5iv.labs.arxiv.org/html/2301.08906v1
"Consistency models in distributed systems: A survey on definitions, disciplines, challenges and applications" - https://ar5iv.labs.arxiv.org/html/1902.03305
When you're talking about say, high frequency trading systems, this is a big deal.
And anyone who deploys a service with all their data in a single data center, is a pretty small time player, or they don't need the uptime.
CAP Twelve Years Later: How the “Rules” Have Changed (2012) - https://news.ycombinator.com/item?id=10179628 - Sept 2015 (4 comments)
CAP Twelve Years Later: How the "Rules" Have Changed - https://news.ycombinator.com/item?id=4043260 - May 2012 (3 comments)
That PutIfAbsent Api would solve so many problems for metadata formats used in the BigData field like Delta or Iceberg!
About two years ago AWS made a big announcement that they could finally guarantee "read your own writes" consistency model. But ONLY within a single account.
If you know that you are not racing against eventual consistency, you can use HeadObject API call to check whether a given key exists in the bucket or not.
Yeah, over the last 5 years but mostly before that change, I saw some strange intermittent problems with Hive where certain files would be registered in the Hive metadata by one process but not visible when the data actually went to get queried by other processes which led to jobs we were responsible for bombing out with missing file errors.
What was truly strange was when this situation would last for hours.
Most services, short of financial ones, choose to sacrifice consistency for availabilty. Example being social networks, it doesn't matter if I get "inconsistent" social media feeds depending on where I am because my European friends posts haven't had time to propogate across the database shard.
OTOH financial systems, you best believe they are going to choose consistency 9 times out of 10. Exceptions being being very small amounts. I don't want someone to exploit network/node downtime in order to steal millions of dollars.
But the reality is that network outages, rarely last very long. Then its just a matter of making sure the nodes themselves are reliable.
It's just recently that anybody started to care about instantaneus coerence, instead of eventual.
Did you mean coherence? As in consistency?
I'd have autocorrected in reading except that you wrote it twice which implies intention. However a quick search didn't turn up a definition.
Of course you could always write bad checks - but payee was on the hook for the fraud, not the bank ...
Also, ironically enough, most of the global financial system actually operates very inconsistently. It can take days for transactions to clear, and there's a bazillion quirks based on local law and bank policy. So in practice banks use double entry accounting with reconciliation. If something got messed up they issue a new compensating transaction.
One hand went up in the back. That's when I knew our culture was doomed.
CAP is (for the record) why the blockchain is and is not "just a database" -- it's a database which takes a different position on CAP and as a result can be used for different things to other databases.
1) there is a network
2) on each block, a leader is selected to make the next block on the basis of a lottery -- each lottery ticket is a unit of computing done ("proof of work")
3) if the network splits, each subnet has a different leader and a different sequence of blocks from that point onwards
4) when the network comes back together, a new "agreed history" is formed from the last shared block, selecting the block with the most "work" done
and resulting in the transactions which were done and confirmed in the smaller subnet being thrown out and history rewritten as if they had never existed
This is the genius of Satoshi -- he figured out that something this trashy was good enough to get something done. Nobody in an academic setting would ever have considered that as a worthwhile solution I don't think.
It was madness. But it worked. It's the same kind of genius as Ethernet just rebroadcasting packets until something works rather than doing Token Ring type negotiation about who can talk next.
Worse really is better at times.
The thing most people pick in practice is CA, single master with hopefully a replicated backup.
That's something lost on a lot of modern engineers, I think. Distributed Systems are something you do because you have to, not because they're the best approach inherently.
The theorem is about being perfectly C, A, and P
If you have more than two hosts, you have P.
If you have delay to sync, you have C.
If you send inaccurate data instead of waiting you have A.
Theres no 'mostly'
so there's a linear relationship between reducing sync time and reducing response time until at some point in reducing sync time you reduce response time below a level you give a shit about...
and you have all 3 of consistency, availability, and partitioning
and then you have beaten the CAP theory
Because you can have a well partitioned system with good data redundancy, and you can take time to let the system become consistent which reduces availability.
But the actual practical outcome of changing availability is changing response time.
Therefore, consistency time is directly and linearly related to response time.
And if the response time is below some threshold that you no longer care about then you have beaten CAP theorem.
So in a practical sense, you can beat the CAP theorem. In an abstract sense you cant.
But response time grounds the theorem to reality because it in every system, there's a point of optimization that no longer matters.
You'd think that would be blindingly obvious, but based on every system in existence ... It's not.
Primarily I think it is marketing/sales since thinking properly about distributed systems is hard and executives don't want to hear about it.
https://arxiv.org/abs/1509.05393
TL;DR: terms in CAP are not adequately defined and he proposes a delay-sensitive framework.
The "C" and "A" in CAP theorem have very specific definitions in the context of the theorem, that isn't the same as, say, the "C" in ACID. The CAP theorem in its original formulation makes no claim about latency, or eventual consistency, or any of these things. It is a very trivial logical conclusion.
"Consistency" just means when a node is queried, you either get the latest read, or you get an error.
"Availability" just means when a node is queried, you always get something. You never get an error.
---
Given this, we can show with a simple example that in the face of a network partition, it is impossible to get C and A simultaneously.
Suppose there's a Leader node and a Follower node with the two being kept in sync, but say there's a network partition between the two, so the Follower cannot reach the Leader.
Say I changed my name in the Leader from "old-name" to "new-name".
If I ask the Follower node what my name is, what should it return?
The follower can't know that the update has occurred, so you can see that it can only do one of two things
- The follower could return an error (it has C, but not A)
- The follower could return "old-name" (it has A, but not C)