5 ms·
Indeed, applying the CAP theorem to real-world databases makes no sense, because the CAP definition of "available" is unnecessarily restrictive. The real trade
by devit 11y ago
Indeed, applying the CAP theorem to real-world databases makes no sense, because the CAP definition of "available" is unnecessarily restrictive.
The real tradeoff is much simpler: if you want a consistent system, it will be slower and more expensive.
Regarding availability in a consistent system, as long as more than half of the servers are working and connected with each other, they will be able to elect a leader and function in a consistent way (using the Raft protocol for instance). Now as long as a client can connect to at least 1/k of the servers in the majority it can just keep trying connecting to random servers until it finds a reachable one in the functioning majority in time proportional to k.
For systems within a single datacenter, redundant networking makes partitions almost impossible, and having enough hosts makes losing more than half almost impossible, so the only failure mode is losing the whole datacenter. For systems distributed among multiple datacenters, availability is almost guaranteed unless a global catastrophe causes a global Internet partition (such as Eurasia and America no longer being connected).
The only issues, again, are that it's more expensive because you need more nodes and redundant networking and that it's slower because nodes have to communicate among each other before commits can be confirmed, especially if that needs to be done across datacenters. In particular, write throughput does not necessarily scale with more nodes in a consistent system, since all writes could modify the same value (after reading it) and that's not parallelizable in the general case.
- seiji 11y agoFor systems within a single datacenter, redundant networking makes partitions almost impossible, derp, nope. No amount of redundant networking can avoid administrator mistakes or accidental "our redundancy all runs through the same conduit and ninjas chopped it in half during a new buildout." better researched anecdotes: https://aphyr.com/posts/288-the-network-is-reliable https://aphyr.com/posts/288-the-network-is-reliable enough hosts makes losing more than half almost impossible Losing network connectivity is indistinguishable from losing hosts though, and losing hosts is indistinguishable from your application just not responding any longer (dumb programming language GC pauses? kernel panic and halted? infinite docker deploy loop? SSD slowly failing and now writes take 1,000,000 times longer than before? Y2038 bug? All applications run on the same timer or compaction schedule and decide to blip for the exact same 30 seconds?). So, the failure mode for hosts is your network failure probability and your host failure probability, and host failure probability includes software failures.
- devit 11y ago> "our redundancy all runs through the same conduit and ninjas chopped it in half during a new buildout." You can just consider that as a failure of the whole datacenter. The ninjas could just as well chop all the external network connections, which would result in actual datacenter failure (from a client/service PoV), so it shouldn't increase the rate that much.
- seiji 11y agoYou don't always have a choice: http://farm3.static.flickr.com/2316/2216487046_b7ca640f56_o.jpg http://farm3.static.flickr.com/2316/2216487046_b7ca640f56_o.... Plus, we live in cloud la la land these days. You have no idea how any of your machines/VMs are connected together. We can assume nothing.
- rdtsc 11y agoThat's not the same thing. A network partition is not equivalent to the the whole datacenter being off nice and clean. That would be nice actually, because it is an and easily testable failure mode. Network partition due to misconfiguration will happen and they have verious interesting corner cases -- multiple plartitions or say partitions between servers but not between clients (clients see the servers, server don't see each other). You can of course say "it will never happen" and just let chance decide what happens do the data in case when a partition happen.
- devit 11y agoNo, the system needs to be properly designed to be consistent, so the data will always be fine (as long as you don't permanently lose all the servers and backups). Partitions and datacenter failures only determine whether the system is up or not, and thus its availability properties.
- hueving 11y ago>derp, nope. Your post was full of good information. why did you have to start it by being an asshole?
- bcoughlan 11y agoA lot of distributed systems literature seems to be too abstract for real world systems engineering. Where does one go to learn about architecting a distributed system on a that works in the real world?
- hueving 11y agoGoogle's papers on F1 and spanner are decent. Unfortunately not many companies do this and a lot of it is tied up as 'proprietary information'.
- adamnemecek 11y agoThe author of this paper, Martin Kleppmann, is writing a book "Designing Data-Intensive Applications" (http://shop.oreilly.com/product/0636920032175.do?cmp=af-strata-books-videos-product_cj_9781491903094_%25zp http://shop.oreilly.com/product/0636920032175.do?cmp=af-stra...). I've been reading it via the O'Reilly immediate access and I think that it's the book you are looking for.
- dwenzek 11y agoIndeed, he written a less formal post on the matter : "Please stop calling databases CP or AP" (http://martin.kleppmann.com/2015/05/11/please-stop-calling-databases-cp-or-ap.html http://martin.kleppmann.com/2015/05/11/please-stop-calling-d...)
- jmtulloss 11y agoI know you said a lot of it is too abstract, but I really enjoyed the readings from UIUC's CS525 when I was in school (years ago). The selected papers are a great "who's who" of distributed systems research. https://courses.engr.illinois.edu/cs525/sp2015/sched.htm https://courses.engr.illinois.edu/cs525/sp2015/sched.htm
- phpnode 11y agoBasho produce a lot of good documentation for Riak which I've found invaluable when creating my own distributed database.
- 11y ago
- brianpgordon 11y agoPartitions being "almost impossible" would be fine for running Facebook or HN. Not so much if you're a payment processor or investment brokerage. If you want your distributed system to be consistent then your application code needs to be able to handle partitions, end of story. The CAP theorem hasn't changed.
- deleted 11y ago[deleted]
- david927 11y agoFrom what I understand, even ATM's are eventually consistent.
- vidarh 11y agoMost things in the banking system is eventually consistent.
- vidarh 11y agoFor a payment processor at least, there's basically nothing that can't be eventually consistent. Anything interfacing with banking has traditionally had a very high tolerance for eventual consistency as many processes have had wildly varying and long settlement timelines. Basically it's been rare for you to be able to guarantee you have a consistent view of an account at any given time.