Let's consign CAP to the cabinet of curiosities
brooker.co.za
brooker.co.za
A bank: No! If region A goes down, do not process updates in B until A is back up! We’d rather be down than wrong!
A web forum: Yes! We can reconcile later when A comes back up. Until then keep serving traffic!
CAP theorem doesn’t let you treat the cloud as a magic infinite availability box. You still have to design your system to pick the appropriate behavior when something breaks. No one without deep insight into your business needs can decide for you, either. You’re on the hook for choosing.
Do you accept writes in us-west-2a knowing the ones in 2b can't see them? Do you serve reads in 2b knowing they might be showing old information? Or do you shut down 2b altogether and limp along at half capacity in only 2a? What if the problem is that 2a becomes inaccessible to the Internet so that you can no longer use a load balancer to route requests to it? What if the writes in 2a hadn't fully replicated to 2b before the partition?
You can probably answer those for any given business scenario, but the point is that you have to decide them and you can't outsource it to RDS. Some use cases prioritize availability above all else, using eventual consistency to work out the details once connectivity's restored. Others demand consistency above all else, and would rather be down than risk giving out wrong answers. No cloud host can decide what's write for you. The CAP theorem is extremely freaking relevant to anyone using more than one AZ or one region.
(I know you weren't claiming otherwise, just taking a chance to say why cross-AZ still has the same issues as cross-region, as I have heard well meaning people say they were somehow different. AWS does a great job of keeping things running well to the point that it's news when they don't. Things still happen though. Same for Azure, GCP, and any other cloud offering. However flawless their design and execution, if a volcano erupts in the wrong place, there's gonna be a partition.)
Financial institutions set up in this century (like paypal or fintecs) usually do want consistency at all cost. Your risk exposure is a lot lower if you know people don't spend or withdraw more than they have.
I think there's also less liability for the banks if they accept and build around the risk of both eventual consistency and fraud. Even now wiring is necessary to pull off a digital bank theft, and even then you need to be extremely fast to get away with the money, and extremely fast to not go to prison, and liquidate it extremely quickly (likely into bitcoin). Even small delays in the transactions would make this basically impossible with modern banking, let alone a full clearing house day.
The only way around this would be to replace stateful apis with some kind of function that takes a reference to a verifiable ledger, that can mutate that and anyone can inspect themselves to see if the operation has been completed. Oh wait :)
Also, in the real world, we can solve problems at a different level - like legal or product.
Example: some ATMs are configured that if the bank computer goes down (or the connection to it), they’ll still let customers withdraw limited sums of money in “offline” mode, keeping a local record of the transaction, that will then be uploaded back to the central bank computers when the connection comes back up.
But, consider this hypothetical scenario: bank computer has a major outage. ATM goes into offline mode. Customer goes to ATM to withdraw $100. ATM grants withdrawal and records card number, sum and timestamp internally (potentially even in two different forms-on disk and a paper-based backup). Customer leaves building. 30 minutes later, bank computer is still down, when a terrorist attack destroys the building and the ATM inside. The customer got $100 from the bank, but all record of it has been destroyed. Yet, this is such an unlikely scenario, and the sums involved are small enough (in the big picture), that banks don’t do anything to mitigate against it - the mitigation costs more than the expected value of the risk
https://cw33.com/news/heres-how-banks-rearrange-transactions...
You've completely elided the quorum thing to make it sound like TFA is nuts.
I don't think TFA is nuts. I do think the premise that engineers using cloud can ignore CAP theorem is wrong, though. It's a decision you must consider when you're designing a system. As long as everything is running well, it doesn't matter. You've got to make sure that the behavior when something breaks is what you want it to be. And you can't outsource that decision. It can't be worked around.
You don't because their clients won't also be isolated from the other servers. That's TFA's point, that in a cloud environment network partitions that affect clients generally deny them access to all servers (something you can't do anything about), and network partitions between servers only isolate some servers and none of the clients.
You can engineer cloud services to be resilient. AWS has done a great job of that in my personal experience. But worst cases can still happen.
If you have an algorithm that runs in n*log(n) most of the time but 2^n sometimes, it can still be super useful. You have to prepare for the bad case though. If you say it’s been running great so far and we don’t worry about it, fine, but that doesn’t mean it can’t still blow up when it’s most inconvenient. It only means it hasn’t yet.
My goodness. So many commenters here today are completely ignoring TFA and giving lessons on CAP -- lessons that are correct, but which completely ignore that TFA is really arguing that clouds move the needle hard in one direction when considering the CAP trade-offs, to the point that "the CAP theorem is irrelevant for cloud systems" (which was this post's original title, though not TFA's).
What does this mean? What is the cloud being down? The cloud isn't a single thing; it's lots of computers in lots of data centres.
You've misunderstood the problem.
You cannot be consistent if you cannot resolve conflicts, taking writes after a network partition means you are operating on stale data.
Honestly the closest we've ever gotten to solving CAP theorum is spanner[0], which trades off availability by being hugely less performant for writes. I'm aware spanner is used a lot in google cloud (under the hood), but you won't solve CAP by having hundreds of PGSQL replicas, because those aren't using Spanner.
In fact, you can test it out, ElasticSearch orchestrates a huge number of lucene databases, a small install will have a dozen or so replicas and partitions, but you're welcome to crank it to as many nodes as you want then split off a zone, and see what happens.
I'm becoming annoyed at the level of ignorance on display in this thread, so I'm sorry for the curt tone, you can't abstract your way to solving it, it's a fundamental limitation of data.
[0]: https://static.googleusercontent.com/media/research.google.c...
Cloud is always available? I have a bridge to sell you.
The cloud isn't what's being described as available. The service is available or not. The service is hosted on multiple computers that talk to each other. If they stop talking to each other, the service needs to either become unavailable (choosing consistency) or risk having not up to date data (choosing availability).
In the same spirit, some businesses are fine with a single db, making backups every night and losing some data when an issue happens. It doesn't mean that people maintaining these systems get to tell others that distributed systems are essentially a non-problem.
You have to consider the system as a whole including external clients since the definition of the availability of the entire system ultimately depends on the point of view of the end users of the system.
The CrowdStrike incident has demonstrated that no matter what you do it is possible for any distributed system to be partitioned in such a way that you have to choose between consistency and availability and if the system is important then that choice is important.
Partition intolerance means you can only proceed when all nodes have been reached. This means partition intolerance not only deals with availability in the yes or no sense, but also in terms of latency.
The author of the article doesn't actually understand the CAP theorem.
> DNS, multi-cast, or some other mechanism directs them towards a healthy load balancer on the healthy side of the partition
Incidentally that's where CAP makes it's appearance and bites your ass.
No amount of VRRP, UCARP wishful thinking can guarantee a conclusion on what partition is "correct" in presence of a network partition between load balancer nodes.
Also, who determines where to point the DNS? A single point of failure VPS? Or perhaps a group of distributed machines voting? Yeah.
You still need to perform the analysis. It's just that some cloud providers offer the distributed voting clusters as a service and take care of the DNS and load balancer switchover for you.
And that's still not enough, because you might not want to allow stragglers write to orphan databases before the whole network fencing kicks in.
So a pretty simple application design can deal with all of this, and that’s the world we live in, and people deal with delays and move on with life, and they might bellyache about delays for the purpose of getting discounts from their vendors but they don’t really want to switch vendors. If you’re a bank in the highly competitive business of taking 2.95% fees (that should be 0.1%) on transactions, maybe this stuff matters. But like many things in life, like drug prices, that opportunity only exists in the US, it isn’t really relevant anywhere else in the world, and it’s certainly not a math problem or intrinsic as you’re making it sound. That’s just the mythology the Stripe and banks people have kind of labeled on top of their shit, which was Mongo at the end of the day.
A good engineer knows that all real-time systems are turn-based, and the network is _always_ partitioned
The whole point of the CAP theorem is that you can't have a one size fits all design. Network partitions exist, and it's not always a full partition where you have two (or more) distinct groups of interconnected nodes. Sometimes everybody can talk to everybody, except node A can't reach node B. It can even be unidirectional, where node A can't send to B, but node B can send to A. That's just how it is --- stuff breaks.
There's designs that highlight consistency in the face of partitions; if you must have the most recent write, you need a majority read or a source of truth --- if you can't contact that source of truth or a majority of copies in a reasonable time (or at all). And you can't confirm a write unless you can contact the source of truth or do a majority write. As a consequence, you're more likely to hit situations where you reach a healthy frontend, but can't do work because it can't reach a healthy backend; otoh, that's something you should plan for anyway.
There's designs that highlight availability in the face of partitions; if you must have a read and you can reach a node with the data, go for it. This can be extended to writes too, if you trust a node to persist data locally, you can use a journal entry system rather than a state update system, and reconcile when the partition ends. You'll lose the data if the node's storage fails before the partition ends, of course. And you may end up reconciling into an undesirable state; in banking, it's common to pick availability --- you want customers to be able to use their cards even when the bank systems are offline for periodic maintenance / close of day / close of month processing, or unexpected network issues --- but when the systems come back online, some customers will have total transaction amounts that would have been denied if the system was online. Or you can do a last state update wins too --- sometimes that's appropriate.
Of course, there's the underlying horror of all distributed systems. Information takes time to be delivered, so there's no way for a node to know the current status of anything; all information from remote nodes is delayed. This means a node never knows if it can contact another node, it only knows if it was recently able to or not. This also means unless you do some form of locking read, it's not unreasonable for the value you've read to be out of date before you receive it.
Then there's even further horrors in that even a single node is actually a distributed system. Although there's significantly less likelyhood of a network partition between cpu cores on a single chip.
I am really saying that I don't buy into the mythology that someone at AWS or Stripe or whatever knows more about the theoretical stuff than anyone else. It's a cloud of farts. Nobody needs these really detailed explanations. People do them because they're intellectually edifying and self-aggrandizing, not because they're useful.
You run a vacation home exchange/booking site, which you run in 3 instances in America, Europe and Asia because base you found that people are happier when the site is snappy without those pesky 50ms delays.
Now suppose it's the middle of the night in one of those 3 regions, so no big deal, but the rollout of new version brings down the database that's holding the bookings for two hours before it's fixed. Yeah, it was just the simplest solution to have just a one database with the bookings, because you wouldn't have to worry about all that pesky CAP stuff.
But people then start to ask if it would be possible to book from regions that are not experiencing the outage while double-bookings would still not happen. So you give it a consideration and then you figure out a brilliant, easy solution. You just
Reddit: 1.31 s
Amazon: 1.51 s
Google: 577 ms
CNN: 671 ms
Facebook: 857 ms
Logging into Facebook: 8.37 s
As long as you have TLS termination close to the end user, and your proxy maintains connections to the backend so that extra round-trips aren't needed, the amount of time large popular sites take to load suggests that people wildly overstate how much anyone cares about a fraction of a blink of an eye. A lot of these sites have a ~20 ms time to connect for me. If it were so important, I'd expect to see page loads on popular sites take <100 ms (a lot of that being TCP+TLS handshake), not 800 ms or 8 seconds.
You can produce a document that says these load times matter - Google famously observed that SERP load times have a huge impact on a thing they were measuring. You will never convince those analytics people that the pesky 50ms doesn’t matter, because they operate at a scale where they could probably observe a way that it does.
The database people are usually sincere. But are the retail people? If you work for AWS, you’re working for a retailer. Like so what if saving 50ms makes more money off shopping addicts? The AWS people will never litigate this. But hey, they don’t want their kids using Juul right? They don’t want their kids buying Robux. Part of the mythology I hate about these AWS and Stripe posts is, “The only valid real world application of my abstract math is the real world application that pays my salary.” “The only application that we should have a strictly technical conversation about is my application.”
Nobody would care about CAP at Amazon if it didn’t extract more money from shopping addicts! AWS doesn’t exist without shopping addicts!
So the observation is that apparently Amazon does not think 50 ms is very important. If they did, their page could be loading about 5-10x faster. Likewise with e.g. reddit; I don't know if that site has ever managed to load a page in under 1 s. New reddit is even worse. At one point new reddit was so slow that I could tap the URL bar on my phone, scroll to the left, and change www to old in less time than it took for the page to load. In that context, I find people talking about globally distributed systems/"data at the edge" to save 50 ms to be rather comical.
Partition isn’t one ethernet cable going bad. It’s all the ethernet cables going bad. Redundant network providers for your data center to handle the idiot with the backhoe isn’t surviving partition it’s preventing it in the first place.
Something like that happened just this year to 13 African nations [1]
1: https://www.techopedia.com/news/when-africa-lost-internet-ex...
It's most common that you have a full partition when nobody can talk to node A, because it's actually offline. And sometimes you've got a dead uplink for a switch with a couple of nodes behind it that can talk to each other.
But partial partitions are really common too. If Node A can talk (bidirectionally) to B and C, but B and C can each only talk to A, you could do something with clever routing tricks. If you have a large enough network, you see this kind of thing all the time, and it's tempting to consider routing tricks. IMHO, it's a lot more realistic to just have a fast feedback between monitoring the network and the people who can go out and clean the ends on the fiber on the switches and routers and what not. The internet as exists is the product of smart people doing routing tricks on a global basis; it's not universally optimal --- you could do better in specific cases all the time, but it's really easy to identify these as a human when something is broken; actually doing probing for host routing is a big exercise for typically marginal benefits. Magic host routing won't help when a backhoe (or directional drill or boat anchor) has determined your redundant fiber paths are all physically present in the same hole though.
I've seen similar things happen in other software under load as well (which can easily cause cascading failures, too)
We lost a preprod database that held up dev and QA for half a day. To bring the server room up to code they moved the outlets, and when they went to move the power cords, someone thought they could get away with relying on the UPS power during the jiggery. Turned out the UPS that was reporting all green only actually had about 8 seconds of standby power in it, which was a little longer than the IT person needed to untangle the cables.
So in the end we were just fucked on a different day than we would have been if we had waited for the next hiccup in the city power.
I mention this because if you've split your cluster between three racks and someone manages to power down an entire rack, you're dangerously close to not having quorum and definitely close to not having sufficient throughput for both consumer traffic and rebuilding/resyncing/resilvering.
It's a slow motion rock slide that is a cluster that was not sized correctly for user traffic plus recovery traffic as the entire thing limps toward a full on outage when just one more machine needs to be rebooted, because, for instance, your team decided that restart a machine every n/48 hours, to deal with a slow leak is just fine since you're using a consensus protocol that will hide the reboots. Maybe rock slide is the wrong word. It's a train wreck. Once it starts you can't stop it, you can only watch and fret.
In my experience it's not. The cloud will glitch. The load balancing algo will break subtly for your workload. Your traffic will get blackholed for no apparent reason. I spent a week trying to convince a cloud provider they fucked up (the time it took for them to give us someone who could run the appropriate tcpdump) once. There was no global outage.
It's on you to determine if this is important for you or not in your case, but you will need to mitigate it above a certain SLA threshold requirement. It's far for consigning CAP to history or a curiosity, which is like saying you will never have network issues if you use "serverless" stuff.
Not talking about anyone in particular, but sometimes I feel people building and using the cloud reach hubris level of surety in their systems - I worked both sides of the fence, and I know the long tails of fuckups that impact customers though...
In addition, sometimes you land on a bad piece of hardware but it's not bad enough to trigger the provider's monitoring. Few months ago we had a bad EC2 instance whose network would drop every 2 hours causing a bunch of random errors before recovering for a little while.
It was a very old version of ES, and the specific behavior that led to the problem has been fixed for a long time now. But still, the fact that something like this can happen in a cloud deployment demonstrates that this article's advice rests on an egregiously simplistic perspective on the possible failure modes of distributed systems.
In particular, the major premise that intermittent connectivity is only a problem on internetworks is just plain wrong. Hubs and switches flake out. Loose wires get jiggled. Subnetworks get congested.
And if you're on the cloud, nobody even tries to pretend that they'll tell you when server and equipment maintenance is going to happen.
When there’s maintenance going on in the server room and a machine they promised not to touch starts having intermittent networking problems, it’s probably a shitty cable getting jostled. There is an entire generation of devs now that have had physical hardware abstracted away and aren’t learning the lessons.
Though my favorite stories are the cubicle neighbor who taps their foot leading to packet loss.
That's been my model since then. I've never worked with a cable so expensive or irreplaceable that I'd want to take a chance on it a second time. I'm sure they exist. I haven't been around one.
CAP or no CAP, chaos will reign.
I think FLP (https://groups.csail.mit.edu/tds/papers/Lynch/jacm85.pdf) is better way to think about systems.
I think CAP is not as relevant in the cloud because the complexity is so high that nobody even knows what is going on, so the just C part, regardless of the other letters, is ridiculously difficult even on a single computer. A book can be written just to explain write(2)'s surprise attacks.
So you think you have guarantees whatever the designers said they have AP or CP, and yet.. the impossible will happen twice a day (and 3 times at night when its your on-call).
The single machine is a beastly distributed system in of itself, with multiple cores, CPUs, NUMA nodes, tiered caches, RAM, and disks. But when it comes down to two writers fighting over a row, it's going to consult some particular physical place in hardware for the lock.
For those unaware, military networks deal with disruptions, disconnections, intermittent connectivity, and low-bandwidth (DDIL).
https://www.usni.org/magazines/proceedings/sponsored/ddil-en...
Just sprinkle the magic "cloud" powder on your system and ignore all the theory.
https://ferd.ca/beating-the-cap-theorem-checklist.html
Let's see, let's pick some checkboxes.
(x) you pushed the actual problem to another layer of the system
(x) you're actually building an AP system
(x) your solution requires a central authority that cannot be unavailable
I'm amazed the author doesn't consider that the load-balancers in this situation are effectively acting as non-voting witnesses and that the whole system hinges on the lbs being able to correctly determine the healthy nodes. In a single-master setup the network partition that cuts the line between the lbs and the current master (but nothing else) is a fun one. Everything would be fine if you could promote a new master but who's gonna tell em'?
I won't plagiarize his text, instead the chapter references his blogpost, "Please stop calling databases CP or AP": https://martin.kleppmann.com/2015/05/11/please-stop-calling-...
(*): rebuttal I think is the wrong word, but I couldn't think of better.
I liked the linearizable explanation, like when Alice knows the tournament outcome but Bob doesn't. A super extension to this would underscore how important this is, and the danger of Alice knowing the outcome but Bob not knowing, just extend the website a bit and make it a sports gambling website. A system under such a parition would allow Alice to make a "sure thing" bet against Bob, so the constraint should be that when Bob's view is stale, it does not take bets. But how does Bobs view know its stale? It has to query Alices view! Lots of mind games to play!
The author seems to not understand what the meaning of the P in CAP
If you treat a partitioned node as "failed", then CAP does not apply. You've simply left it out cold with no read / write capability because you've designated a "quorum" as the in-group.
Sure. Also, there's a long list of other things that are probably irrelevant to you. That is, until your provider fails and you need to understand the situation in order to provide a workaround.
And slapping "load-balancers" everywhere on your schema is not really a solution, because load-balancers themselves are a distributed system with a state and are subject to CAP, as presented in the schema.
> DNS, multi-cast, or some other mechanism directs them towards a healthy load balancer on the healthy side of the partition.
"Somehow, something somewhere will fix my shit hopefully". Also, as a sidenote, a few friends would angrily shake their "it's always DNS" cup reading this.
edit: reading the rest of the blog and author's bio, I'm unsure whether the author is genuinely mistaken, or whether they're advertising their employer's product.
What a convenient world where the client is not affected by the network partition.
[0]: https://en.wikipedia.org/wiki/Somebody_else's_problem#Dougla...
Iow: You can have CAP as long as you can communicate across "partitions".
Let's say you have servers in DCs A and B, and clients at ISPs C and D. Normally servers at A and B can communicate, and clients at both C and D can each reach servers at A and B.
If connectivity between DCs A and B goes down there is technically no partition, because connectivity between the clients and each DC is still working. The servers at one of the DCs can just say "go away, I lost my quorum", and the clients will connect to the other DC. This is the most likely scenario, and in practice the only one you're really interested in solving.
If connectivity between ISP C and DC A goes down, there is no partition because they can just connect to DC B.
If all connectivity for ISP C goes down, that is Not Your Problem.
If all connectivity for DC A goes down, it doesn't matter because it no longer has clients writing to it, and they've all reconnected to DC B.
To have a partition, you'd need to lose the link between DCs A and B, and lose the link from ISP C to DC B, and lose the link from ISP D to DC A. This leaves the clients from ISP C connected only to DC A, and the clients from ISP D connected only to DC B - while at the same time being unable to reconcile between DCs A and B. Is this theoretically possible? Yes, of course. Does this happen in practice - let alone in a close-enough split that you could reasonably be expected to serve both sides? No, not really.
DCA: north america
DCB: south america
ISP C: north american ISP
ISP D: south american ISP
A major fiber link between north and south america is cut, disrupting all connectivity between the two areas. "lose the link between DCs A and B": check.
"lose the link from ISP D to DC A.": check.
replace north and south america with any two disparate locations.In the early days that might have two machines we walked up to seeing different results, now we talk to everything over the internet. We generally mean partition when one vlan is having trouble. One ethernet card is dying or has a bad cable. If the entire thing becomes unreachable we just call it an outage.
It's simply the formalization of a fact, and whether or not that fact is *important* (although still a fact) depends on the actual use case. Hell, it applies even to services within the same memory space, although obviously the probability of losing any of the three is orders of magnitude less than on a network.
Can we please move on?
Doubtful. If there's one thing that new developers have always insisted upon doing it's telling the networking and data store engineers that their foundational limitations are wrong. It seems that everyone has to learn that computers aren't magic through misunderstanding and experience.
They’re digital! It’s ones and zeroes! Dude, it’s a fucking magnet. A bunch of them in fact.
One of the fundamental assumptions of CAP theorem is that you can't tell whether or not you have a partition. If you have an oracle that can instantaneously tell you the state of every subsystem, then yeah, CAP is pointless.
But if one of your DBs is connected, reporting itself as alive, and throwing all its writes into /dev/null, you won't be able to route traffic to a quorum of healthy instances because it's not possible to be certain that they're all healthy.
This is what CAP theorem is about: managing data in a distributed system where the status of any given system is fundamentally unknowable because of the Two Generals' Problem (https://en.wikipedia.org/wiki/Two_Generals'_Problem)
In many cases in Cloud though, we can skip that technical stuff and design systems as if we really _did_ have an oracle that could instantaneously and perfectly tell us the state of the system, and things will typically work fine.
Unfortunately the point is lost because of the usage of the word "cloud", a somewhat contrived example of solving problems by reconfiguring load balancers (in the real world certain outages might not let you reconfigure!), and missing empathy that you can't tell people not to care about how the semantics that thinking about, or not thinking about, availability imposes on the correctness of their applications.
As for the usage of the word cloud: I don't know when a set of machines becomes a cloud. Is it the APIs for management? Or when you have two or more implementations of consensus running on the set of machines?
Or he's saying you don't need Consistency because your system isn't actually distributed; it's just a centralized system with hot backups.
It's unclear what he's trying to say.
No idea why he wrote the blog post. It doesn't increase my confidence in the engineering equality of his employer AWS
This is the key, that network partitions either keep some clients from accessing any servers, or they keep some servers from talking to each other. The former case is uninteresting because nothing can be done server-side about it. The latter is interesting and we can fix it with load balancers.
This conflicts with the picture painted earlier in TFA where the unhappy client is somehow stuck with the unhappy server, but let's consider that just didactic.
We can also not use load balancers but have the clients talk to all the servers they can reach, when we trust the clients to behave correctly. Some architectures do this, like Lustre, which is why I mention it.
I see several comments here that seem to take TFA as saying that distributed consensus algorithms/protocols are not needed, but TFA does not say that. TFA says you can have consistency, availability, and partition tolerance because network partitions between servers typically don't extend to clients, and you can have enough servers to maintain quorum for all clients (if a quorum is not available it's as if the whole cloud is down, then it's not available to any clients). That is a very reasonable assertion, IMO.
It's not "partition tolerance" if you declare that partitions typically don't happen.
I agree that in modern data centers the CAP theorem is essentially irrelevant for intra-DC services, due the uptime and redundancy of networking H/W (making a partition less likely than other systemic failures).
Across DCs I'll claim it is still absolutely relevant.
Most databases don't work like Spanner, and Spanner has its downsides, two of them being cost and performance. So most of the time, you're using a traditional DB with maybe a RW replica, which will sacrifice significant consistency or availability depending on whether you choose sync or async mode. And you're back to worrying about CAP.
But the gist, I guess, is that for most applications it’s not actually that important, and that’s probably true. But when it is important, “the cloud” is not going to save you.
As always Kleppmann has a great and deep answer for this.
https://martin.kleppmann.com/2015/05/11/please-stop-calling-...
However, in the example with the network partition, it relies on proper monitoring to work out if the DB its attached to is currently in partition.
managing reads is a piece of piss, mostly. Its when you need propagate write to the rest of the DB system, thats where stuff gets hairy.
Now, most places can run from a single DB, especially as disks are fucking fast now. so CAP is never really that much of a problem. However when you go multi-region, thats when it gets interesting.
What if your servers can't talk to each other, but clients can?
What if clients can't connect to any of your servers?
What if there are multiple partitons, and none of them have a quorum?
Also, changing the routing isn't instantaneous, so you will have some period of unavailability between when the partition happens, and when the client is redirected to the partition with the quorum.
Then it’s not a partition.
I suppose it's kinda true in the sense that how to operate a power plant is not relevant when I turn on my lights.
You still get consequences for ignoring them, but they show up as burn rate.
Umm, no? That’s a picture of a partition. The partition is not able to make progress because the system is not partition tolerant. If it did it wouldn’t be consistent. It’s still available.
> In database theory, 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).
I likely need to read the paper linked, but it's common to have an MPP database lose a node but maintain data availability. CAP applies at various levels, but the notion of availability differs:
1. all nodes available 2. all data available
Redundancy can make #2 a lot more common than #1.
At that point you get all 3: consistency,availability, partitioning.
In my opinion it should be the CAPR theorem.
That being said, if this is truly a problem for you CRDB is basically built with this all in mind.
That's an incredible piece of engineering.
I don't think it makes sense to say that CAP doesn't apply if you don't need consistency, availability, or tolerance to partitions. CAP is entirely about the need to relax at least one of those three to shore up the others.
I would like to violate CAP, please. I would like to be nearish to c, please.
Here is my passport. I have done the work.
Yeah… no. Just because the cloud offers primitives that allow you to skip many of the challenges that the CAP theorem outlines, doesn’t mean it’s not a critical step to learning about and building novel distributed systems.
I think the author is confusing systems practitioners with distributed systems researchers.
I agree in some part, the former rarely needs to think about CAP for the majority of B2B cloud SaaS. For the latter, it seems entirely incorrect to skip CAP theorem fundamentals in one’s education.
tl;dr — just because Kubernetes (et al.) make building distributed systems easier, it doesn’t mean you should avoid the CAP theorem in teaching or disregard it altogether.