Two Phase Commit
martinfowler.com
martinfowler.com
"Do you take X? I do. And do you take Y? I do. I now pronounce you X+Y".
If the coordinator crashes, one (or all) of the transactions will be stuck until the coordinator comes back - or, if it never does - until someone manually rolls them all back of commits them all.
If your definition of "consistency" is "read will always returns the latest value which is accepted by all nodes in the first phase" you'll always have to contact all nodes when you're reading the data (since as you mentioned you can't actually depend on if it's been "committed"). This is "CP" because you "tolerate" a network partitioning by having the read return an error if it can't reach all nodes.
But you can "tweak" your definition of "consistency" to be "reads will read the latest committed value on that node". CAP still applies for the system as a whole (with writes), but now you can serve fast reads from a single node.
Many systems defines "read consistency" as "snapshot isolation" which means (roughly) that all requests inside a single transaction will be executed against a single snapshot (which possibly isn't the "latest" one). This will also allow you to execute reads against a single node.
And to quote the two authors who actually proved CAP, Seth Gilbert and Nancy Lynch. (Emphasis mine):
> Consistency, informally, simply means that each server returns the right response to each request, i.e., a response that is correct according to the desired service specification. (There may, of course, be multiple possible correct responses.) The meaning of consistency depends on the service.
> Please don't post shallow dismissals, especially of other people's work. A good critical comment teaches us something.
The main thing I'm disagreeing with is undefining what CAP means. CAP is a very particular proof and the rest I think perhaps we need a new name for.
> CAP has a long history and was originally provided without a proof.
Irrelevant. All of database history is filled with what today is comparatively junk. The understanding has improved and we have better ways of talking about the particular subject matter.
> The proved version of CAP has a very strict definition, but it's not that useful in practice because most systems don't actually follow that model.
It is useful in what it proves, so don't bend the definition. This is my main point. You will have to invent another taxonomy perhaps. We have good names for all the typical consistency models I believe, so we can use those as a starter.
> Most notably, the CAP theorem is often still true when "consistency" is defined as weaker than "linearizable".
This is the type of reasoning that just muddies rather than clarifies.
Raft/Paxos is to make sure majority of the replica for a partition have the latest version of the data and that you can keep writing even if minority of nodes in that partition are unavailable. while still giving guarantee that no write get lost and all read reflect the most recent write.
While 2PC is to write 2 records atomically when the 2 record are located on 2 different partitions.
edit: Looking at the paper, it looks like the tradeoff is increased amount of coordination needed.
"The Two-Phase Commit protocol is thus the degenerate case of the Paxos Commit algorithm with a single acceptor."
I suppose in Spanner, having more than one acceptor is redudant since each shard is a Paxos group anyways.
If the coordinator goes down, the system makes no progress.
If the coordinator sends inconsistent commands (e.g., commit to one resource and abort to another), it's not serving its purpose.
So in a system using two phase commit, your availability is limited by the availability of the coordinator.
If the coordinator is a single server, you have a single point of failure in a distributed system.
If you try using multiple servers to increase the availability of the coordinator, you risk sending inconsistent commands -- unless you implement a consensus system within the coordinator.
Paxos and Raft let you combine the availability of multiple servers while still letting them behave consistently. (This is the "consensus problem", and correctly solving it is notoriously tricky).
EDIT: found it! https://lamport.azurewebsites.net/video/consensus-on-transac...
i am being a bit hand-wavey here because what the paper is talking about is an application of paxos to generate a new algorithm (paxos commit) that devolves to 2pc.
my approach has been to make all operations idempotent and ensure they are all ran at least once.
1. User clicks "buy now" for whatever is in their shopping cart 2. Client generates some kind of transaction ID representing that they wanted to purchase the contents of the shopping cart (could be deterministic ID) 3. Client submits this request to the server 4. Server persists the intention to start processing the purchasing of the shopping cart with transaction ID of X 5. Server synchronously or asynchronously starts handling the side effects of the purchase 6. If at some point the client got an error message it can still submit the same request with the same transaction ID to retry and even if the initial request was received (but perhaps lost before getting to the client) it's cheap and easy to make it idempotent by using the transaction ID
Race conditions would be made more difficult by having everything idempotent based on the transaction ID and having the transaction ID (optionally) generated deterministically.
2 phase commit is an extremely heavy weight pattern and finds far less use than something like the above.
> generated deterministically
in this case generated deterministically means generated by some immutable value based on the initial transaction properties (who is buying what, with x quantity at y time) and not just a random uuid?That being said, I don't think it's a good idea for reasons similar to this: https://wiki.c2.com/?DistributedTransactionsAreEvil vs. using the Saga pattern.