4 ms·
To be clear, my sense is that every replica would receive every transaction, and would be free to commit the transactions inside each batch in any order. Now I
by quinthar 10y ago
To be clear, my sense is that every replica would receive every transaction, and would be free to commit the transactions inside each batch in any order.
Now I do agree that it's tricky to avoid "gaps" in the failure cases. However, every replica keeps a record of the past several million transactions (we aim for 3 days), and every transaction is assigned a unique ID. When a replica starts up, it "synchronizes" down every missing transaction it has, and at this point would "repair" any gaps it somehow obtained when it went down.
Admittedly, the exact details of that part are TBD, but it doesn't strike me as an unresolvable problem on the surface.
- Jweb_Guru 10y agoThe question is: (1) In the event of a crash, how does it know which transactions it's missing without a total order on the write transactions? (2) In the non-failure case, how do you know the transactions in one batch are ordered before the transactions in a subsequent batch (which requires the replica to identify any "gaps" within the previous batch so it can wait for them to come in before continuing on)? I don't think any of this is unresolvable (just adding a total order would go a long way). I do think it's very tricky to get right, with lots of edge cases, and that if you can't efficiently identify missing transactions it potentially makes recovery unusably slow. In any case, since I believe you are losing serializability with multiple writers, you should probably weigh that against any performance gains you get from multiple writers (I think even the group commit variant suffers from this problem).
- Jweb_Guru 10y agoAlso, as pointed out in https://twitter.com/aphyr/status/788757992829222912 https://twitter.com/aphyr/status/788757992829222912, even your current solution isn't safe if your commit order isn't total across leader changes (you can resolve this by adding an epoch number that increments every time leader election occurs). Strongly recommend you take him up on his offer.
- quinthar 10y agoTo be clear, we're talking about functionality that is not implemented or fully designed. Today all transactions are committed on all nodes in the same order, which is a much simpler world. I agree, the multi-threaded replication case is a much more complex and interesting world, with much greater performance opportunities. Lots of exciting problems to solve when we get there!
- Jweb_Guru 10y ago> Today all transactions are committed on all nodes in the same order, which is a much simpler world. This is difficult to reconcile with: > - For the highest performance, you can designate a transaction as "asynchronous" and the master will commit immediately because if the leader crashes, a replica becomes leader and starts accepting writes, then the old leader recovers as a replica, without something like an epoch number it won't be able to tell that it has commits that the current leader doesn't (using a unique incrementing transaction number based on just a counter at the leader won't work, because it won't necessarily be unique across leader elections thanks to the asynchronous commits).
- quinthar 10y agoAh, sorry for the confusion. Every transaction is given an incrementing ID by the leader, and every follower commits the transactions in ID order. Furthermore, every commit has a running SHA hash of all prior commits (and every node keeps a history of the last few million commits). This way any two nodes can compare their journals to make sure they agree -- and if there is any split, then the cluster kicks that node out. Basically, there is no scenario in which a node that commits a different transaction (or a transaction in a different order) is allowed to remain in the cluster.
- Jweb_Guru 10y agoI think this or something like this can probably work if you're okay with losing all data that wasn't acked by a majority (though I suspect actually recovering a divergent replica would be very difficult), but this doesn't work with the batch commit idea at all, does it? Seems like it enforces strict serial ordering of writes (even nonconflicting ones).