602 karma · joined August 13, 2008
The linked-through article is much better.
edit: more precise wording
Let's imagine a system where you want to be 100% available for reads, for any number of failures less than N. Then you need to be able to submit every single write to every single node in the system, otherwise the failure of all but the up-to-date node will result in stale reads.
But then if a single node is partitioned from the network, we can't (correctly) be available for writes, because the system is incapable of sending updates to all reads as required. It doesn't matter which node you ask.
The point is that every system has a failure mode like this. I take your point that it's not always just a single node failure that precipitates the abandonment of C or A, but that was never the point of the CAP theorem.
I don't agree - availability is a totally meaningless property if you are allowed to occasionally return "no, I won't process your request". Such a response communicates nothing about the state of the atomic object you are writing to or reading from, so you can always return it and trivially satisfy 'availability' if we define it this way.
To your other point - be aware that I didn't write the article, so I'm not speaking for the author. However, I think you're right that the article makes it sound a bit like a single failure or message loss will cause any protocol to immediately sacrifice availability or consistency. This isn't the case - all CAP does is establish that for every protocol, there exists a failure pattern that will force it to abandon one of the two.
For quorum systems, this means that a permanent partition causes one half of the partition to no longer be the majority, and therefore can no longer be consistent if it responds to any requests. Paxos is another example.
So you're right, there are particular patterns of partition that stop a protocol from functioning correctly. And many that don't - hence the term 'fault tolerant' has some meaning.
Avoiding these patterns, in practice, can turn out to be surprisingly tricky. High-performance systems can't afford to have too many participants, which means that the probability of a problematic failure is higher than we might like (five participants in a consensus protocol is already a lot for high throughput, but now we are susceptible to only three failures). Failures are also often correlated, so independence assumptions don't hold as much as would like. Machines crash. Networking gear fails.
There's no abuse of statistics here. The probability of a particular failure pattern can be engineered low, and at that point you must weigh the trade-offs of the cost of loss of availability / consistency vs. the effort you make to minimise the chance of occurrence. We are talking about edge cases here, and the implicit assumption is that the cost of hitting one of them is huge (and it often is). However if you run a cluster large enough, you hit edge cases all the time.
(Although you mention partial synchrony, note that most of these results are mainly applicable to asynchronous networks in the first instance).
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.
The theory behind all this really does hold this point up. I have another blog post with much more detail on the theory here: http://the-paper-trail.org/blog/?p=49, but I warn you it may be heavy going.
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.
Paxos is, fundamentally, a quorum-based system that deals with reordering of messages. It sacrifices liveness for correctness - if the proposer does not hear back from a majority of nodes (in the case of, e.g. a partition), the protocol will not complete (availability is sacrificed).
My point is not that there is a 'vital packet' in every protocol, the omission of which will cause either a lack of availability or consistency (although I can certainly design protocols that way!) - it's that for every protocol there is a network partition which causes it to be either unavailable or inconsistent. That network partition might be dropping ten messages, or just one. Retransmitting would make sense, but in real life message failures are often highly temporally correlated :(
The proof of this, by the way, is in a very famous paper by Fischer, Lynch and Patterson called "The Impossibility of Distributed Consensus With One Faulty Process". One take away is that one slow-running process can take down any protocol. It may take a few missed messages, but only a single node...
http://www.cloudera.com/blog/2010/04/cap-confusion-problems-...
Therefore it's not a CA system, but a C system.
Reporting an error condition counts as an availability violation.
If the network is allowed to drop packets there are times where either you must not respond to requests (as doing so would violate sequential consistency) or you must respond incorrectly, potentially with stale information.
The network partition that forces you to drop one of these guarantees might be quite dramatic - but note that, for example, a quorum system is only available on the majority side of its partition. Therefore if a single node is unable to deliver messages (due to a network partition event), it will not be able to correctly respond to requests and must either not respond to its clients or respond incorrectly.
The consistency guarantee requires that RW histories are compatible with some sequentially consistent history on a non-concurrent RW register. Defining a total order on operations is sufficient, I believe, but not necessary (does it matter what order two consecutive reads happened in?).
UX, UI, PM, distributed systems engineer, operations engineer and more. We're in the Bay Area, down in Palo Alto and are genuinely a great company for which to work.
If you're interested in any of our positions, drop me a line at henry at cloudera.com and I'll get you in contact with the relevant people - in particular if you're a distributed systems guy looking for some seriously interesting problems to work on, I'd love to hear from you!
Threads are a reasonable structuring tool for expressing concurrency, which is useful for laying out your code in a maintainable, easy-to-reason-about way.
Mainstream Python implementations have concurrency but really do not have parallelism.
In order to preserve the illusion of independence, the OS has to deal with the possibility that the sum of the sizes of all the memory that each process wants to use might be greater than the amount of physical memory available. So rather than aggressively limit the amount of virtual address space that each process can use, it simply only keeps a subset of that memory in physical memory at any one time.
You can have virtual memory without paging, but then each process has to compete for a very limited resource. You can also have paging without virtual memory - process A's copy of physical address X can be swapped out to be replaced with process B's copy. However processes are then still limited by having only as large an address space as physical memory in the machine, and virtual memory is such a huge win for hiding the layout of physical memory from processes as well as isolating them from one another (so no corruption possibilities) that it's pretty much unheard of to do this.
You get such a big advantage from having a layer of indirection sometimes...
A trivial commit protocol that is consistent but not live simply sends no messages. All updates that succeed are consistent, but no updates succeed.
An eventually consistent protocol is often correct for stronger consistency guarantees, but instead sacrifices consistency in the case of network problems rather than liveness.
(sidebar: this is really what the CAP theorem is driving at. You must choose between consistency and liveness if the network can lose messages).
Vector clocks, at the cost of more storage, allow you to detect whenever two messages are not causally related. To be concrete, if I send you a message, and you then send one to someone else, your message is causally related to mine because its possible that whatever was in my message influenced the message you sent. However, if I send my message to a third party, it no longer causally precedes yours because you couldn't have reacted to what I sent.
Vector clocks will preserve this lack of ordering because the clock of my message will be 10, and the clock of your message will be 01. 10 is neither 'less than' or 'greater than' 01 under the ordering of vector timestamps, so we can conclude that the messages are not causally related.
This is important because sometimes we care whether two different updates to the same data are in any way ordered - because if so we know which one was 'latest' and therefore the 'current' one. Dynamo uses vector clocks like this: if two non causally related message arrive to update a value, the system is able to detect a conflict and take some remedial action.
Computer science shares some of the creative sensibilities of mathematics: the building blocks (theorems) may be well understood, but there is a creativity and insight required to arrange them in such a way as to produce an [efficient/elegant/small/general] method to solve a given problem (proof). At the same time, this is why coding on its own is such a pale shadow of the entire field of computer science. Implementations at their worst are merely transcriptions of someone else's work.
Don't aspire to be an artist with your code; it's not the medium for art. There's a different aesthetics at work here. You will not habitually provide commentary on society or the human condition. Don't begrudge others the fact that they might. They don't begrudge you the beauty that you produce.
The most often quoted example of Dynamo's use is the shopping cart application on every Amazon page. In the worst case, your shopping cart will mysteriously empty itself. This is a huge pain, and a potential loss for Amazon, but it's not catastrophic in the way that is implied here. Indeed, assuming the liveness of a quorum, the application will read back all conflicting entries for the shopping cart (those that aren't ordered under their vector clock timestamps) and the onus is on it to merge the conflicts. Of course, the shopping cart will take the union of all updates to ensure that nothing is dropped (and therefore some delete operations may be lost).
The key point is that some applications can do without observing a linearisable history, and the interest of this paper is that it explores the design space if you drop that requirement.
I don't understand the post's points about CAP; all three requirements are in tension. Dynamo is unusual in that it is live in the case of a network partition while still maintaining its consistency guarantees.
Similarly - those systems that use chain-replication asynchronously like he describes can still suffer from the same read-old-value-after-it-was-written consistency issue, if the reader jumps between two replicas for consecutive reads. Avoiding that can require synchronous coordination of updates (a la Paxos, e.g.) which is, I think, what the paper is driving at. Otherwise, there are still failure modes which, in order to patch up, require stronger guarantees about liveness of quorums than Dynamo needs.
I understand that Dynamo is no longer used internally at Amazon at scale, so maybe some of the practical points this post makes about the realities of central coordination held water for real deployments. Still, I don't buy the reaction that prioritising availability uber alles and designing a system that does not behave exactly like a strongly-consistent key-value store immediately invalidates it for workloads that have high availability requirements and lower consistency needs.
http://www.enterpriseintegrationpatterns.com/docs/IEEE_Softw...
I like Lamport's opinion of his own contribution:
"My other contribution to this paper was getting it written. Writing is hard work, and without the threat of perishing, researchers outside academia generally do less publishing than their colleagues at universities. I wrote an initial draft, which displeased Shostak so much that he completely rewrote it to produce the final version."
This is one of two great fundamental results in consensus protocols; the other being the later FLP result which shows the impossibility of consensus with failures in an asynchronous network. Byzantine fault tolerance was a hot topic in systems for the last roughly ten years thanks to Turing Award winner Barbara Liskov's work with Miguel Castro on Practical BFT that sparked a bit of a resurgence; the two most notable other papers IMHO in the field since that are 'Separating Agreement From Execution' paper from UT Austin and the SOSP '07 paper on Zyzzyva. Both of these take apart the 3f+1 bound and show exactly how many processors you need for various parts of the protocol.
----
You’re doing a good thing by trying to elucidate basic concepts, but I’m afraid your article has a bunch of errors that means it’s not as helpful as it could be.
“Big O specifically describes the worst-case scenario” - this isn’t true. g(n) = O(f(n)) simply means that f(n) asymptotically bounds g(n) as n->inf - but the functions involved could be anything and don’t necessarily refer to the worst case. For example, they could be the average case running times of an algorithm, or something completely different. Also remember that for any n, g(n) could be > f(n), so f(n) is not worst case in that respect either.
“O(1) describes an algorithm that will always execute in the same time (or space) regardless of the size of the input data set.” - curiously, also not true. It just means that there’s eventually an upper limit on how much space an algorithm will use, or that eventually as n gets larger the resource consumption of the algorithm becomes constant.
“The example below also demonstrates how Big O favours the worst-case performance scenario; a matching string could be found during any iteration of the for loop and the function would return early” - same point about g(n) = O(f(n)) not being worst case necessarily - you’ve set the problem up so that N is the worst (and average) case running time. However we can say that the best case running time of linear search is O(1).
“O(2^N) denotes an algorithm whose growth will double with each additional element in the input data set.” - also you need to remember that these are upper bounds, and aren’t necessarily tight. That linear search example, which you wrote as O(N), is also O(2^N). If we replace O by \Theta then we begin to get to the behaviour you’re describing, but remember that these things only are true as N gets sufficiently large.
Yes, you're correct that I actually don't get as far as the computation. That's coming next - at over 2,000 words it became apparent that the series needed to be segmented.
(I should have called this part 0, or maybe part aleph-0 :))
I won't be blazing new trails here though. My aim is to cover, roughly, the following:
* Turing's attack on the Entschiedungsproblem and the halting problem (if space permits, some context regarding Godel) * The correspondence between Turing Machines and natural numbers, and the corollary that most real numbers are uncomputable. * Rice's theorem, and recursive and recursively enumerable sets. * Possibly some mention of Chaitin's Omega.
All should be covered in an undergraduate course on computability - however I see such misunderstanding of simple ideas like the halting problem that come from a shaky grasp of the unintuitive basics that I wanted to write a genuine introduction.
I'd appreciate any further suggestions for content!