CAP Theorem Explained
robertgreiner.com
robertgreiner.com
Once you are distributed, P is not an optional. Rather, in the case of network failure, consistency or availability is what suffers. A system cannot be both CA and distributed.
So, for example, Elasticsearch is not a 'CA' solution despite the diagram in the article, but is actually closer to PC (although, in practice it is far more subtle than even that as it is not perfectly consistent, and configuration options allow for some trade-offs between availability and consistency in the case of communication errors.
Thanks for your feedback though, I have some posts queued up on more of the subtleties in the model, this was just meant to be an introductory post.
If you configure ES to prioritize consistency somewhat (minimum_master_nodes), it prevents writes during a partition—but there's at least one "split" partition scenario where even minimum_master_nodes doesn't prevent inconsistent writes. If you configure ES to prioritize availability during a partition, it isn't consistent. Remember, ES doesn't claim to be consistent and doesn't even use any sort of consensus algorithm.
Re: the debate about CA vs CP, I've seen this too, but I think it reflects confusion rather than genuine options. The 'P' of CAP takes effect whenever you have two or more processes which clients can communicate with. So even if you are running within the same datacenter when using ES, provided clients can communicate with more than one ES node, CAP still applies.
Longer version (approx. four years old) here:
http://pl.atyp.us/wordpress/index.php/2010/10/when-partition...
I think I also might write one up myself.
http://blog.foundationdb.com/minimal-explanation-of-the-cap-...
But, what I think is actually true is: Partition tolerant means that nodes still respond to requests while some other nodes are out of reach (one node being one or more instances of the DBMS in one cluster).
Am I right?
So, If that is true, having a RDBMS with master<=>master replication between two sites is Available and Partition tolerant, albeit inconsistent (the replication is not possible).
In the case of a master<=>master replication as a load-balancing/failover solution, it's about availability and consistency, but we should not even call that "distributed".
Still right?
Geography and whether it's a "failover solution" doesn't change that the system is distributed - anything with more than one node is fundamentally distributed, and I'd argue that even two processes communicating on one instance can be considered as a distributed system as well.
The system can choose which kind of predictable it wants to be: whether it sacrifices consistency in favor of being available (accepting reads and writes without knowing that reads are fresh or writes linearizable) or availability in favor of consistency (rejecting reads and writes that can't be confirmed).
In the case of "master<=>master replication" as it's configured in most relational databases, I think the system tries to be AP: conflicted writes or stale reads are possible because barring special options like 2-phase commit, the replicas generally lag each other (as anyone who's tried to reconcile a MySQL split-brain situation knows, this can be a tremendous pain).
* http://research.microsoft.com/apps/pubs/default.aspx?id=1926... - including a very nice way of thinking about CAP and tradeoffs)
* http://www.infoq.com/articles/cap-twelve-years-later-how-the... - a good perspective on CAP, and why many people still don't understand it clearly.
* http://lpd.epfl.ch/sgilbert/pubs/BrewersConjecture-SigAct.pd... - don't get put off by the formal language. Gilbert and Lynch is still (IMO) the best explanation of what CAP means, and what it implies.
* http://blog.cloudera.com/blog/2010/04/cap-confusion-problems... - Some good criticism of the way CAP is frequently explained.
* http://codahale.com/you-cant-sacrifice-partition-tolerance/ - Why CA systems don't actually exist.
* http://cs-www.cs.yale.edu/homes/dna/papers/abadi-pacelc.pdf - PACELC, maybe a better model.
* http://dbmsmusings.blogspot.com/2010/04/problems-with-cap-an... - Another look at PACELC, with some examples.
This post has some info on elastic, but nothing in excruciating detail: http://aphyr.com/posts/288-the-network-is-reliable
"YOU CANT REDUCE FUNDING TO PARTITION TOLERANCE. YOULL REGRET THIS."
I wonder if anybody's made SimDataCenter?
That would be really interesting.
Therefore, I am curious in knowing. Here's my question: how do they configure production systems in order to get reasonable or 100% maybe consistency? In other words, the reads must factor latest write into account. No, stale reads.
I thought that consistency in ACID meant that data was always consistent with the rules of the database, whereas in CAP it means that the same data held in different locations is the same? Is that not right?
Consistent means that they are consistent from the perspective of the client within the rules of the consistency model. If you have sequential consistency and quorum operations, your replicas do not need to have the same data, but the view of your clients is always consistent.
From Eric Brewer himself: "The relationship between CAP and ACID is more complex and often misunderstood, in part because the C and A in ACID represent different concepts than the same letters in CAP..."
(from http://www.infoq.com/articles/cap-twelve-years-later-how-the...)