4 ms·
Every part of the cap theorem is commonly misunderstood. Availability is probably the most commonly misunderstood aspect as is outlined in this article. People
by jeffffff 9y ago
Every part of the cap theorem is commonly misunderstood. Availability is probably the most commonly misunderstood aspect as is outlined in this article. People commonly conflate consistency with durability when really CAP says nothing about durability because CAP is based on a model where the only type of failure is a network partition. And of course there are the people who claim network partitions don't happen in their network so they can choose CA.
CAP is also not stated with a sharded database in mind. Availability is stated as follows:
"For a distributed system to be continuously available, every request received
by a non-failing node in the system must result in a response"
This means that if you have 100 machines and you store each row on 3 of those, if the 3 machines storing a row are partitioned from the client then the client can read and write the row on any machine. That machine must respond with something like "the row doesn't exist" or "the row is empty" if queried about the contents. This quickly devolves into an unusable system if you follow this train of thought to its conclusion. No practical sharded database is truly AP under this definition and it only makes sense to talk about CAP in the context of a single shard of a sharded database (which can be AP).
It is actually quite easy to have consistency without sacrificing latency: Don't replicate. What you are sacrificing in this scenario isn't consistency, it's durability, which is totally ignored by CAP. Where the latency is really introduced is in trying to achieve durability. A durable AP system will end up with close to the same write latency as an equivalently durable CP system, but no one bothers to build AP systems with the same level of durability as a CP system. It is possible to build a durable CP system where reads are just as fast as in any AP system as long as there are no partitions at the time of the read. The trick is to use read leases so that reads only have to touch one machine.
The CAP theorem is nowhere near as important as many people think it is with regards to performance and uptime of real world systems. It is a theoretical result based on a very simplified model.
The only thing AP buys you is handling certain types of network partitions that almost never happen as you said. They have better uptime when you can talk to some machines but less than a quorum. This is not common.
The real reason AP systems are common is because consistency+durability is really really really hard.
- canes123456 9y agoIf you are not replicating, than it not a distributed system and cap doesn't apply. Sharding just gives you two centralized databases. Talking about durability is just adding more confusion IMO.