3 ms·
I believe Fauna inherits some scalability issues from Calvin, in that given NxN communication is necessary for transactions to be processed, a catastrophic fail
by _benedict 4y ago
I believe Fauna inherits some scalability issues from Calvin, in that given NxN communication is necessary for transactions to be processed, a catastrophic failure in one shard can bring the entire database offline. Is that correct?
It may not happen in practice today, but as your clusters grow your exposure to such a failure is increased, and the NxN communication overhead grows does it not?
It's certainly more scalable than other systems offering this isolation today (besides Spanner), but algorithmically at least I believe the proposal we are developing for Cassandra is strictly more scalable, as the cost and failure exposure grow only with the transaction scope, not the cluster.
Not to ding FaunaDB, though. It probably is the most scalable database offering this level of isolation that is deployable on your own hardware - assuming it is? I know it is primarily a managed service, like Spanner.
I'm also not aware how large any real world clusters have gotten with FaunaDB in practice as yet. Do you have any data on that?
- bitbckt 4y ago> I believe Fauna inherits some scalability issues from Calvin, in that given NxN communication is necessary for transactions to be processed, a catastrophic failure in one shard can bring the entire database offline. Is that correct? No, that's not correct. FaunaDB is inspired by Calvin, but is not a direct implementation of Calvin. Transactions are applied independently on each storage host, though checking OCC locks may involve peers within a single replica. > assuming it is We used to offer on-premises delivery, but we are strictly a managed service now. > I'm also not aware how large any real world clusters have gotten with FaunaDB in practice as yet. Do you have any data on that? I'm not at liberty to say.
- _benedict 4y ago> Transactions are applied independently on each storage host I’m not talking about transaction application, but obtaining your slot in the global transaction log. How does a replica know it isn’t missing a transaction from some other shard if that shard is offline come its turn to declare transactions in the log? It must at least receive a message saying no transactions involving it were declared by that shard, no? This is pretty core to the Calvin approach, unless I misunderstand it. > I'm not at liberty to say. I think scalability is something that is a function of both theoretical expectations and practical demonstration.
- bitbckt 4y ago> How does a replica know it isn’t missing a transaction from some other shard if that shard is offline come its turn to declare transactions in the log? It must at least receive a message saying no transactions involving it were declared by that shard, no? A transaction is committed to the log with sufficient information necessary for the log (and each storage host) to know whether any particular host is involved in that transaction. There's no need for a message to prove a host's absence from that transaction - it has all it needs in the log to determine that it isn't involved. There is some complexity around how that information maps to the cluster's physical topology, but hosts which aren't involved in a transaction don't process that transaction and need not check with any other host to know that fact. [ETA] That information is derived from the read set when constructing the transaction to propose to the log. Logical partitions that weren't a part of that read set don't need to know that they weren't involved - that fact is obvious from their absence. > I think scalability is something that is a function of both theoretical expectations and practical demonstration. I agree with you, though as a commercial product operated as a service, that practical demonstration is naturally limited in this forum.
- _benedict 4y agoYou misunderstand my point. A replica may of course safely assume it is not involved in a transaction it hasn’t witnessed so long as it has a “complete” log (as in, the portion involving it). > that fact is obvious from their absence. Yes, but absence from what? There must be some message that contains the relevant portion of the log declared by each other shard. If that shard is offline, so it cannot declare its transactions at all, how does the system continue? This absence is one step back from the one you are discussing. If a shard’s leader is offline, a new leader must be elected before its slot comes around for processing - and until this happens all transactions in later slots must wait, as it might have declared a transaction that interferes with those a later slot would declare. If no leader can be elected (because a majority of replicas are offline) then the entire log stops, no?
- bitbckt 4y agoIt's true: I've had trouble determining what exactly your point is, as you seem to think you know more about FaunaDB's architecture than you do. You're under the mistaken impression that we have the same data/log architecture as described in the Calvin paper, which requires that every host talk with every other host. That is not true of FaunaDB.