3 ms·
Great insight from Abadi as always. To add to his point about guarantees and network partitions: I find that the way consistency is explained in documentation o
by thamer 3y ago
Great insight from Abadi as always. To add to his point about guarantees and network partitions: I find that the way consistency is explained in documentation or even courses can be somewhat misleading or at least incomplete. If you read the Dynamo paper and take from it that you could model a system with a "classic" design of 3 replicas doing quorum reads and writes, you could easily convince yourself that it will be consistent and continue operating even with the loss of a replica.
What this simple approach doesn't cover is all the cases that aren't your perfectly well-behaved read and write operations. It doesn't cover the case where your client gets a timeout and has no idea how many replicas were written to. You could have one replica that has persisted the (failed) write, and read back an old copy from the other two. Then your next read could pick up the new version, giving you alternating views of this data that was supposed to be consistent.
With a database like Cassandra you also have to consider that only the most recent copy of a value is returned, regardless of how many replicas responded. In that case you could read an old value, then the new one, then the old one again, all at quorum when it was supposed to be consistent.
Failures are often messy, and the source of most of the complexity in reasoning about consistency in distributed systems. Many concepts (like network partitions) are also often misunderstood, or being considered in a way that's too "theoretical". Nodes don't always shut down instantly. You likely won't see a clean network split that starts at a fixed point in time and is also resolved in an instant. It'll be partial, asymmetric, with nodes coming in and out, etc. Good luck reasoning about these scenarios…
- EdwardDiego 3y ago> Nodes don't always shut down instantly. You likely won't see a clean network split that starts at a fixed point in time and is also resolved in an instant. It'll be partial, asymmetric, with nodes coming in and out, etc. Excellent point, it's precisely this scenario that has caused the biggest issues in distributed systems I've been working with.
- YZF 3y agoWith Cassandra if you wrote at quorum and it was successful and you read at quorum you will get your value back. Where it gets hairy is when a write fails (e.g. only one replica got it) or if writes so close to each other that the clocks are a problem. It is definitely non-trivial to build systems over that but then it's non-trivial to build systems that scale horizontally, and geographically, and are highly available, and can recover from failure, in the first place. Like the article implies, what tends to happen is people don't understand the guarantees and build software that works except when it doesn't. Then you build on top of that etc.
- _benedict 3y agoIt’s not possible to read a new value at quorum and to then later read the old value at quorum in Cassandra. This is enforced by “read repair” that ensures the quorum agrees on the reply, by propagating the newer values if they are missing on any of the quorum, before replying. That’s not to say failures in distributed systems aren’t messy; they definitely are.