7 ms·
You Can't Sacrifice Partition Tolerance
- rcoder 16y agoI think the yield/harvest concept described here is one of the more useful models I've heard about in a while for thinking about fault-tolerance tradeoffs. My thanks to @codahale for the write-up, and particularly for the references.
- ericflo 16y agoThis is the most clear and well-written summary that I've seen of the tradeoffs presented by the CAP theorem. Hopefully it clears up a lot of the confusion out there. I think this is the tweet that prompted the post: http://twitter.com/JamesMPhillips/status/26502076366 http://twitter.com/JamesMPhillips/status/26502076366
- moshezadka 16y agoI would disagree with the central thesis. You can sacrifice P for a weaker version: assume a network which "eventually heals": any live node will answer at least one message in a hundred, say [any node which does not is assumed to be unavailable]. The alternative to P is not "perfect network", it's "bounded from below on the badness thereof network", a significantly more realistic beast.
- rlpb 16y agoA CA system is simply one which is not available at all during a network partition, since it is partition-intolerant. This lack of availability is different from the availability in the A of CAP, since that availability holds only so long as the network is not partitioned (by definition in a CA system). Such a system might not be considered a distributed system at all (although it may still be distributing load), since a partition-intolerant system is effectively one system as far as the CAP theorem is concerned. So it's essentially a special case of the CAP theorem, but it is still useful to describe it as CA.
- HenryR 16y agoNo, it's exactly the same. Availability is a guarantee that all requests are eventually responded to within some time bound, whatever that is. During the partition, availability is violated. Therefore it's not a CA system, but a C system.
- dadkins 16y agoAre you sure? Availability in the CAP theorem is a state, as are (P)artition and (C)onsistency. Your system can't be simultaneously consistent and available in the presence of a network partition. The A in CAP doesn't mean always available. It just means the system can, at best, be any two of the three at a time.
- HenryR 16y agoNo, it does mean always available - honestly :) If there is some time period during which requests are not responded to within a time bound, the system is not available then, and further is not a 'highly' or 100% available system. That is what the CAP theorem is talking about. Consistency, similarly, is not a state but a property that holds across all responses. Either you return a consistent response to all your requests, or you don't. In the context of CAP, there is no middle ground.
- dadkins 16y agoI think the author, like many recently exposed to the CAP theorem, is confused about the meaning of partition tolerance, leading to ridiculous conclusions. Partition tolerance does not mean your distributed system can't be consistent and available because your network dropped one packet, or one node failed. What would be the point of such a definition? Instead, the CAP theorem implies that while the network is partitioned, consistency or availability must be sacrificed. In the case of the dropped packet, once it is retransmitted the partition is healed and progress can be made. Or in the case of the failed node, nothing says that the rest of the system can't be consistent and available, so that the system as a whole maintains that property. There is no requirement that the unavailable node be available. Truly partition tolerant systems are those that continue to function in the face of a prolonged partition, and those are the systems that must sacrifice either consistency or availability.
- allertonm 16y agoWhat you're saying is that so long as there are no partitions, the system can be Consistent and Available, but if there's a partition, it can't. "Consistent sometimes" is not the same thing as "Consistent" and "Available sometimes" is not the same thing as "Available" - and so "Consistent and Available sometimes" is not the same as "Consistent and Available". I believe you might be guilty of confusing "Eventual Consistency" with "Consistency". Funnily enough, no-one has found much use for "Eventual Availability" so far.
- mcodik 16y agoAssuming eventual availability can be pretty handy -- one way deal with a dependency outage is to retry with an exponential backoff. If a dependency is unavailable now, and your system keeps retrying until it is, then you are assuming your dependency will become available again eventually.
- allertonm 16y agoFair point - but of course your client is not making any progress, and so the unavailability ripples up. It's unlikely that there is user facing case where this is a useful way to work, though I can see it's use in loosely coupled connections between backends.
- jeffffff 16y agothe best way i've heard it phrased is 'given the presence of a network partition, you must choose whether to maintain consistency or availability'. this does not mean that any network partition will make data unavailable if you choose C. it only means that some network partitions will make some data unavailable to some machines. picking A does not guarantee that all data will be available to all machines in the presence of a partition either. given a CP system with 3 way replication requiring a quorum to make progress, i would argue that the set of partitions in which data becomes unavailable yet would still be available had AP been chosen is very small and not worth worrying about. in systems designed to be up 100% of the time where partitions are the exception rather than the norm CP is almost always the right choice. in systems designed for network partitons, like replication to mobile devices or laptops or whatever, AP is almost always the right choice. the problem with trying to apply the CAP theorem to the real world is that the CAP theorem's definition of availability is not the same as most people's definition of availability in practice.
- rlpb 16y ago> in systems designed for network partitons, like replication to mobile devices or laptops or whatever, AP is almost always the right choice. Although it is pretty much the only choice for general purpose file sync, it's still not a good choice. It is difficult to train non-technical staff to deal with inconsistencies on resync, and they don't want to have to deal with it. I've had success with Synctus precisely because it guarantees consistency (it is CP, and any one node keeps the A for a given file). Of course, this only works for mostly-on systems.
- jeffffff 16y agoyes that's certainly true. if there isn't a sensible option for a merge policy and non technical users have to resolve it manually you're fighting a losing battle in an AP system.
- Deestan 16y agoI take this to mean that if the nodes are disconnected from each other, Synctus disallows all access to certain files on node A; i.e. the files that node B currently owns. Do I understand correctly?
- antirez 16y agoI don't agree. For instance Redis Cluster will be consistent (under the limit of physics) and not partition tolerant. But why this requires a network that will never have troubles? Simply when the network will be broken the cluster will not work at all. What Redis Cluster will guarantee is that you can have M-1 nodes, with M being the number of replicas per "hash slot", that can go down, and/or get partitioned. So this is a form of "weak" tolerance to partition, where at least a given percentage of the nodes must remain up and able to talk to each other. But in the practice this is how most networks work. Single computers fail, and Redis Cluster will be still up. Single computers (or up to M-1) can experience networking problems, and Redis will continue to work. In the unlikely condition that the network is split in two halves the cluster will start replying with an error to the clients. This means that the sys admins have to design the network so that it is unlikely that there are strange split patterns, like A and B can talk to C that can tolk to D but blablabal... in high performance network with everything well-cabled and without complex routing this should not be a problem, IMHO.
- inklesspen 16y agoIn your first paragraph's example ("Simply when the network will be broken the cluster will not work at all.") you are sacrificing availability. In the rest of your post, you seem to be sacrificing consistency; one server is down, and thus not receiving any updates from the other servers when data gets updated. I'm not sure you understood the point of the article, so I'll try to restate it: When part of your system goes down (and it will), you can choose between refusing requests, in which case you sacrifice availability, or serving requests, in which case you sacrifice consistency, since the part of the system which is down cannot be updated when you update data, or cannot be queried in the case of data which is insufficiently replicated. You _cannot_ choose both, since that would require communicating with the downed server.
- antirez 16y agowhy do you think my servers are interconnected? I think your conclusions are broken because of many non-always-true assumptions. In Redis Cluster there is no cluster data communication if not for resharding that only works when the whole cluster is on and is done by the sys administrator when adding a node. So in normal conditions, a node will either: 1) Accept a query, or 2) Tell the client: no, ask instead 1.2.3.4:6380 All the nodes are connected only to make sure the state of the cluster is up. If there are too much nodes down from the point of view of a single node it will reply to the client with a cluster error. What I'm sacrificing is only consistency because in every given time there is only a single host that is getting the queries for a given subset of keys. The exception is in the resharding case that is also fault-tolerant. Or slave election (fault tolerance is obtained via replicas). As a side note, the clients should cache what node is responsible for a given set of keys, so after some time and when there are no failures/resharding in act, every client will directly ask the right node, making the solution completely horizontally scalable. Dummy clients will just do always the ask-random-node + retry stage if they are unable to take state. Edit: there are little fields like this that are totally in the hands of academia. My contribution is from the point of view of a dummy hacker that can't understand complex math but that will try to be much more pragmatic.
- lusis 16y agoI guess I'm missing something because the concept of quorum deals with partition tolerance. You require, to provide an answer, that X nodes agree on the state of the data. 3 nodes, 2 must agree. When the partition heals, the resolution process happens. It would have to be a SERIOUSLY bad network design and quorum setting that allows a quorum on both sides of the split. It's just like eventual consistency. We're not talking days or even minutes. We're talking milliseconds/seconds of partition split. If you have a partition split for days, you have OTHER issues to address.
- tbrownaw 16y agoFor a distributed system to be continuously available, every request received by a non-failing node in the system must result in a response. For a CA system, any node which is unable to assure global consistency reports itself as failed. It will neither return bogus results, nor hang indefinitely.
- HenryR 16y agoIn asynchronous networks it is surprisingly hard to detect failures, even of yourself. Reporting an error condition counts as an availability violation.
- tbrownaw 16y agoIn asynchronous networks it is surprisingly hard to detect failures, even of yourself. The reason failures are hard to detect in asynchronous networks is that permissible message transit times are unbounded; ie they refuse to acknowledge the presence of any partition. If you acknowledge the possibility of partitions, then your system is by definition not asynchronous. Reporting an error condition counts as an availability violation. This is bullshit. Per the definition quoted in the linked article, availability only means that "...every request must terminate.". It is not required that it terminate successfully.
- HenryR 16y ago"The reason failures are hard to detect in asynchronous networks is that permissible message transit times are unbounded; ie they refuse to acknowledge the presence of any partition. If you acknowledge the possibility of partitions, then your system is by definition not asynchronous." No. Like you say, async means failures are hard to distinguish from delays. If a node's NIC sets on fire, I'm pretty sure no messages are ever going to get delivered to it - hence it is partitioned from the network. It is very hard to tell whether it has failed, or whether it is just running slowly, in an async network. "This is bullshit. Per the definition quoted in the linked article, availability only means that "...every request must terminate.". It is not required that it terminate successfully." No. The definition of the atomic object modelled by the service doesn't include an 'error' condition. Otherwise I could make a 100% available, 100% consistent system by always returning the error state, which is thoroughly uninteresting. You have to read more than the quoted definition in the Gilbert and Lynch paper to start calling bs - it is very clear that authors do not allow an 'error' response.
- HenryR 16y agoThis blog post that I wrote a few months ago also explains the same issue, and may be of interest for those looking for a separate explanation: http://www.cloudera.com/blog/2010/04/cap-confusion-problems-with-partition-tolerance/ http://www.cloudera.com/blog/2010/04/cap-confusion-problems-...
- ericflo 16y agoAn update: Eric Brewer, who originally posited the CAP theorem, endorses this article: http://twitter.com/eric_brewer/status/26819094612 http://twitter.com/eric_brewer/status/26819094612