3 ms·
You only have to write to the coordination state when there is a failure. You can commit millions of transactions in the happy case without ever doing such a w
by voidmain 8y ago
You only have to write to the coordination state when there is a failure. You can commit millions of transactions in the happy case without ever doing such a write. And failure detector performance and other engineering concerns are usually more of a limitation, in practice, on the performance of recovery than the latency of the coordination state consensus, even when the coordinators are geographically distributed.
- infogulch 8y agoSo the strategy is to optimistically assume that there are no failures and just replicate to all N+1 copies. If there's a failure then back off to the consensus state to coordinate the fix rigorously. In the best case with no failures this works great. But as the number of failures increases, I feel like due to the extra synchronization there will be an inflection point where the cost of the extra layers of coordination will be higher than just synchronizing the data directly. But due to 'other concerns' that inflection point is pushed back by a lot. Is that a reasonable characterization?
- voidmain 8y agoIf you expect to have lots of (hopefully very temporary!) node failures, I think FoundationDB has another trick up its sleeve. You can store (say) N+2 replicas of transaction logs, which are also relatively small and (since sequential) efficient. Then you have a write quorum of N+1 and a recovery quorum of 2 logs, and you don't have to do coordination on every failure. It's certainly true that with enough failures you aren't going to make much progress. I'm not sure that is any less true with plain old state machine replication protocols, though.