684 karma · joined June 13, 2014
Optimal is probably to use linear probing for a few indices at a time (depending on size of hash table entries, say, 4 or 8), and double hashing to move to a new region if no empty elements can be found in that run.
Once we offer interactive transactions we’ll have to detect cycles, but that should be relatively straightforward since execution dependencies are managed explicitly.
Yes, it is very similar in effect, with the absence of a global log conferring some important differences with pluses and minuses. I expect to offer interactive transactions eventually, for instance, which may or may not be difficult for Calvin, and faults on a “log shard” are isolated with Accord, but have global impact with Calvin. But snapshot queries are much easier with Calvin.
So, property (1) provides this for non-overlapping commutative transactions. If w1 writes c1, then finishes; w2 writes c2 then finishes; then if r1,r2,r3, etc don't overlap either w1 or w2 either, then they must correctly witness their execution order and be "externally consistent".
If transaction executions overlap, we must work harder to ensure we see each others' strict serializable order, and so we simply avoid ordering ourselves with transactions that cannot affect our outcome. If these other transactions cannot affect our outcome, then we can't witness anything externally about them, by definition.
If r arrives at C1 before w1, w1 must witness it and take a dependency on its execution. If it arrives at C2 after w2, then r takes a dependency on w2’s execution. In this case, w1 will not execute until r does, and r will not execute until w2 does, and so it is simply not the case that w1 happened first - w1’s commit will not have been acknowledged to the client yet, and the actual order of execution will end up w2,r,w1, which is entirely correct to all observers since w1 and w2 will now overlap their execution.
> operate over all nodes in the database regardless of the shards the keys fall into
Nope, it’s all entirely independent for each shard. The (implied) dependency graph is what ties shards together, so only when they interact.
If w1 and w2 overlap their execution then any serializable execution is strict serializable, so w1 must have finished before w2 began. In this case (2) applies and w1 must have the effect of executing before w2, to all observers, including r.
If w1 and w2 are commutative then it doesn’t matter in what order they execute as they’re equivalent. If not, then at some point their dependency graphs intersect and will be/have been ordered wrt each other by the same restrictions.
edit: in the example given in the Jepsen report, with Accord r would have to execute either before w1; after w1 and before w2; or after them both. This is because one of the following must occur:
1) r takes a dependency on both w1 and w2 and executes after them both;
2) r takes a dependency on only w1; w2 takes a dependency on r, so w2 executes after r, and r executes after w1
3) w1 and w2 both take a dependency on r, so that r executes before either
1) That the result of an operation is reflected in the database before a response is given to the client
2) That every operation that starts after another operation's response was given to a client has the effect of being executed after.
Accord enforces both of these properties. Accord only avoids enforcing an ordering (2) on transactions that are commutative, i.e. where it is impossible to distinguish one order of execution from another. This requires analysis of a transaction and its entire graph of dependencies, i.e. its conflicts, their conflicts, etc. So, if your most recent transaction operates over the entire database, it is ordered with every transaction that has ever gone before.
Cockroach does not do this. If it did, it would be globally strict serializable.
The system maintains a logical clock, and the logical clock enforces the strict serializability condition globally. The logical clock is just cheaper to maintain if the coordinator proposes a value that is already correctly ordered by the loosely synchronized clocks.
edit: sorry, I misread your post as interpreting the loosely synchronized clocks as being necessary for correctness. I answer (very loosely) how strict serializability is maintained in a sibling comment, but in brief a dependency set of all conflicting transactions that might potentially execute earlier is maintained, and a transaction executes in logical timestamp order wrt these dependencies (and, hence, also their dependencies). If a transaction and its dependencies never interacts with another transaction or its dependencies, then their execution order is commutative, i.e. any ordering is indistinguishable from any other.
Nope. That's precisely what Accord manages to maintain. It is not the first such protocol to do so, but so far none of the others have left the lab.
> for scalability, this insertion process happens in batches, and the log itself is replicated and partitioned
This is the link's total exposition on this topic, but it is consistent with my reading of the Calvin paper.
If any log partition becomes unavailable, the log in its entirety becomes unavailable until that partition recovers[1].
Anyway, it feels like we’ve wasted enough of each others’ time without making much progress. Might as well leave it there.
[1] ... and every log partition must communicate with every replica, hence NxN (or MxN)
If not, we’re now in an unfortunate situation where neither your practical nor theoretical capabilities are public knowledge.
> every host talk with every other host
Only every shard
> you seem to think you know more about FaunaDB's architecture than you do
You are perhaps being uncharitable. The only public statements I am aware of about Fauna's capabilities relate to its Calvin heritage. I have only been asking if my understanding is correct regarding both Calvin and its application to Fauna. However, none of the prior issues you responded to were problems with the original Calvin paper either, at least by my reading.
> 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?
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.
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?
Only Spanner currently offers properly scalable strict serializability, as FoundationDB clusters have a size limit (that is very forgiving, and enough for most use cases).
Apache Cassandra is working on providing scalable (without restriction) strict serializable transactions[1], and they should arrive fairly soon. So far as I am aware, at this time it will be the only distributed database besides Spanner offering fully scalable and fast global transactions with this level of isolation.
[1] https://cwiki.apache.org/confluence/display/CASSANDRA/CEP-15... (Disclaimer) I'm one of the authors
I was surprised and impressed with just how well they managed to mimic the vibe despite being live action, but I guess everyone experiences these things differently.