Inconsistent Thoughts on Database Consistency
alexdebrie.com
alexdebrie.com
Let's say transactions are committed to a distributed log (ala raft), then the log position can be the transaction id.
Passing the transaction id back to the client will then allow it to choose its own level of consistency by:
1. Pick consistency over latency: Make read request with transactions id to a random node, and the node wait until it has caught up that point in the distribute log before responding
2. Pick latency over consistency: Make read request for whatever is latest data on random node
Note how this is a departure from how raft services client reads. In raft land only the primary is allowed to service reads as it waits for other consensus tick to ensure it is still the leader before responding. While this guarantees consistency, it hurts scalability as you now have a single bottleneck serving all clients for BOTH read/writes.
I imagine variations of this are quite common, vector clocks as distributed barriers.
But I don't agree that the CAP theorem only applies "when you're in network partition." An available (AP) datastore will have potentially inconsistent data if you write-then-read (ie, the write that you applied may not be read back) even if you "aren't in network partition."
Using that assumption - no partition - why can't we determine if a bit data on a node is stale and request it from somewhere else?
Think through the options here:
1) When one node writes new data, it'll notify the other nodes that it wrote data. Thats called eventual consistency, and you'll get stale reads.. because of the latency between the time the item is written and the time the remote nodes get notified of the update.
2) On write, you make sure all of the nodes accept the update. This is strong consistency, with all of the positives and negatives that includes (high latency, reduced availability).
3) On read, you check with all of the nodes for the newest data. This is another form of strong consistency.
If you weaken the last one, and say: I'll set a timeout, so if there's a network error (even without a partition, there could be a network error) it doesn't block... if that timeout is triggered, then you havent checked all of the nodes, and so you dont know if it's stale (in some cases it will be).
There's literally no way to guarantee the data is not stale without strong consistency... and the laws of nature are the reason why. Data cannot be transmitted instantaneously. Therefore you have to handle that delay somewhere. You can accept stale reads (eventual consistency), or you can accept high latency reads or commits (strong consistency). But the latency never disappears.