9 ms·
CockroachDB is not strict serializable. It is linearizable+serializable+(some other guarantee I forget). Only Spanner currently offers properly scalable strict
by _benedict 4y ago
CockroachDB is not strict serializable. It is linearizable+serializable+(some other guarantee I forget).
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%3A+General+Purpose+Transactions?preview=/188744725/188744736/Accord.pdf https://cwiki.apache.org/confluence/display/CASSANDRA/CEP-15...
(Disclaimer) I'm one of the authors
- richieartoul 4y agoYou’re right about CockroachDB, my mistake. I think your characterization of Cassandra as “fully scalable” in a way that other systems are not is misleading, but I won’t argue it with you :p
- _benedict 4y agoI'm not talking about Cassandra being more generally scalable than other systems, only that no system as scalable as Cassandra offers strict serializable transactions besides Spanner (including Cassandra, today)
- richieartoul 4y agoWhat does scalable mean in this context? In general, Cassandra is much more efficient at writes and can handle higher volumes of write throughput assuming identical hardware, but from what I’ve seen FDB is as efficient (if not more so) for read heavy workloads than Cassandra. A small FDB cluster with no tuning running on small VMs will easily saturate each nodes 10gb NICs for read queries without breaking a sweat. The cluster size limitations are also greatly overstated and mostly a hold over from stale documentation from 2013. You can (and companies do) run FDB clusters with 100s or 1000s of nodes without issue. Cassandra supports this too but in my experience it gets quite difficult to operate at that point, although I’m sure recent versions have gotten much better about this. If you only care about raw write throughout I think your statement is mostly accurate, but for general purpose workloads I would argue that FDB is a system “as scalable as Cassandra” that also offers strictly serializable (and interactive!) transactions which makes it much more broadly applicable to a wide variety of use cases.
- bitbckt 4y agoFaunaDB also offers scalable strict serializability. Disclosure: I work on FaunaDB.
- _benedict 4y agoI 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.
- kbr- 4y agoWhat do you mean by "scalable without restriction"? I understand scalability by being able to handle more traffic through adding more nodes to your cluster. But if you want to handle more traffic, adding more replicas to a single Accord replica set won't help. Any two (fast path) quorums must intersect. And that intersection point for two queries is where you need to spend resources to handle both queries. So - unless I'm missing something here - in order to scale out to be able to handle more traffic, you must shard the data, there's no way around it. Even assuming 0 failures won't help. And once you shard your data and use multiple Paxos/Raft/EPaxos/Accord/whatever groups, you lose strict serializability across the shards.
- _benedict 4y ago> you lose strict serializability across the shards 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.
- kbr- 4y agoOk. From skimming through the paper, it seems that for a given transaction, only the replica sets of the shards touched by this transaction need to participate. Thus for two transactions touching disjoint sets of shards, they can run completely in parallel. Hence the scalability. Is that right?
- _benedict 4y agoPretty much.
- grogers 4y agoIt seems like accord uses neither bounded synchronized clocks nor a shared timestamp oracle. How does it provide strict serializability then?
- evancordell 4y agoI’m confused by this as well, it seems to only provide “strict serializability” for overlapping transactions, which other databases like Cockroach also provide.
- _benedict 4y agoAccord provides global strict serializability. Loosely speaking, strict serializability requires two things: 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.
- richieartoul 4y agoI’m not 100% sure, but my reading of the CockroachDB Jepsen test: https://jepsen.io/analyses/cockroachdb-beta-20160829 https://jepsen.io/analyses/cockroachdb-beta-20160829 indicates that it does meet those two requirements, but that it still is not globally strictly serializable due to the presence of an anomaly they call “causal reversal” where: “transactions on disjoint records are visible out of order.” If you read their Jepsen report and blog posts carefully, the anomaly they suffer from is not caused by a read transaction that started after previous write transactions had completed, but by a read transaction running concurrently with two other write transactions which both commit while the read is still running, and the read sees the writes in the wrong order. The definition you’ve provided above does not cover this case (I think) but its still a violation of strict serializability which makes me think your definition is too loose. Correct me if I’m wrong though, at this level things get quite confusing.