The space between theory and practice in distributed systems
brooker.co.za
brooker.co.za
My suspicion, outside of johnparkerg's polarization, is that distributed systems in practice are particularly messy from a formal standpoint, in contrast to non-distributed systems, while in practice you can live with the messiness, if you can reduce it a sufficient number of 9's, which seems to be anathema to theoretical approaches.
For example, the proof of the CAP theorem is irrelevant in practice, specifically because no system I'm aware of makes the strong consistency assumptions that it requires. On the other hand, CAP behavior is definitely a problem, once you reach a certain size.
"But it's not mathematically stable!"
"It has been for the last month..."
Good question. I think the answer is that you don't build large systems using consensus; you bootstrap large systems with very small systems using consensus. The very small systems are reasonably assumed to have no Byzantine faults (i.e. just like you more or less rely on a single database server not to have faults).
All the systems I know of that use consensus are meant to be small, e.g. run on 5 machines or so (and of course the membership is fixed). Google's Chubby, Yahoo's ZooKeeper, and similar systems like doozer and etcd all work like this.
Consensus doesn't "scale" anyway (the latency isn't bearable). If you only have 5 machines, the likelihood of Byzantine faults over a reasonably long time period is low. The main problem you will see is your own software bugs (i.e. not bugs due to faulty CPU, memory, disks, switches, etc.).
I should add this question to one of my exams and see if my students get it right. I'll be a little grumpy if they don't, but it's a great question. :)
I suppose you can mathematically construct some kind of non-total message loss such that Paxos would never make progress. But that kind of message loss won't persist forever in a real system... it would basically be message loss with knowledge of the algorithm?
A: Prepare(1)
B: Prepare(2)
A: Accept(1, Va) # fails.
A: Prepare(3)
B: Accept(2, Vb) # fails.
Etc... In practice randomized backoffs when this happens causes it to converge extremely fast, but the exact right situation COULD cause it to delay forever. Don't lose sleep over it though.
It is sort of Byzantine in that the "system" is inserting very specific message losses with knowledge of the algorithm.
I believe the true problem lies in our need to categorize people as either scientists and theoretical engineers or down-to-earth engineers, which only polarizes the spectrum.
There is a spectrum of design choices between strong consistency with best effort availability and best effort everything that make sense to many practical use cases.