Don't Settle for Eventual Consistency
queue.acm.org
queue.acm.org
Some business processes/problems can handle eventual consistency. Banking is the classic (and perhaps un-expected) example. You can overdraw your account, there are processes to go back and solve conflicts. Some are automatic, some are manual some are legal.
Bitcoin at the blockchain level is eventually consistent, just wait 10 minutes or whatever the current time it (I don't use it personally just know about it) and then you can be fairly sure of the validity of the transaction. There are built-in incentive to assure histories will converge. But, wallets kept at some exchange should _not_ be eventually consistent. Should _not_ be able to take 100x more than your wallet holds and send it to someone else. There is no regulatory, automatic of any other kind of framework to revert transactions that went through.
There is crdt (a commutative replicated data type) research. So these are data types that an always solve inconsistencies should they arise and instead of diverging they auto-converge, in face of conflicts. Think set union operation or max() function.
These kind of trade-off will percolate up through your data layer into your business problem. For some cases you'd want to pick one, for some pick another.
You don't have to go for full serializability - you can very often get away with something simpler like consistent writes and potentially out of date reads. That sort of system scales a long way unless you're very write heavy.
Part of the reason we can't transfer money instantly between two accounts in the USA is due to sorting out eventual consistency. Federal guidelines on transfer intentionally make it a slow process so the manual processes can catch up.
Yeah, but that's because everything is a case of eventual consistency. Causality itself is limited to the speed of light, and "instantly" is impossible. The only question is whether you want to block/wait, or gloss over it with eventual-tricks.
Just look at online FPS games! Even with some of the best consumer-grade communication links and high expectations for each node, it's impossible to provide actual "instant" behavior, and all modern games contain huge reams of code dedicated to maintaining an eventually-consistent environment.
I think a hybrid consistency model will end up becoming the way we end up going, but not without some changes to the working definition of eventual consistency.
I think that eventually consistent systems that we'll see in the future will at a minimum have atomic, consistent (as in ACID), durable transactions. The big thing you're missing is isolation, and by that, you're also missing serializability. (http://www.vldb.org/pvldb/vol7/p181-bailis.pdf)
Bounded staleness consistency will probably make most people happy -- as in "Give me a consistent snapshot of the database as of 1 second ago," or "Give me a consistent snapshot of the database as of this logical time" (assuming the database consumer has some sort of logical clock you're passing back and forth -- In essence, MVCC with some exposure of the timestamp. (better explanations: http://pages.cs.wisc.edu/~cs739-1/papers/consistencybaseball...)
In addition to this, part of the problems that were outlined in the paper talked about multi-datacenter issues -- a lot of issues with banking-like apps can be solved with escrowing For example, each datacenter has 10% of the account's balance, and you can use consistent commits in the datacenter. If an interactions effects more than 10% of the account's value, it should occur across datacenters, and make a consistent commit across datacenters. (More info: http://mdcc.cs.berkeley.edu/)
Anyways...the future is bright, but we need someone to throw a ridiculous amount of money, and time at it.
And by hybrid in this case I mean some part to keyspace can be tagged as consistent while others remain as before.
Github for the more recent system: https://github.com/wlloyd/eiger
You can find the first author's posters and talks on the subject of scalable casaully-consistent storage here: http://www.cs.princeton.edu/~wlloyd/research.html
The "head in the sand" algorithm isn't always reliable.
They need to understand the tradeoffs, and be able to decide which guarantees are important when.
You can't just use an eventually consistent store as an ACID system and expect it to work, and it's almost always a Really Bad Idea to try to implement ACID on top of EC. Understand how your database works and design to its strengths instead of trying to pretend that its weaknesses don't exist.
This is why I like the idea of Aggregate Roots from domain-driven-design. The boundary is clear in your object model.
"Oh you hit the save button? Yah, we'll get to that."
When a distributed system is partitioned, it is impossible to execute all operations on a shared resource serializably. So, the system must either fail some requests, or fail to be serializable.
You can have your principle of least surprise, or you can have a system that is capable of serving traffic even after some drunk sailor drags an anchor through a submarine cable. You can't have both.