Partitioned consensus and its impact on Spanner’s latency
dbmsmusings.blogspot.com
dbmsmusings.blogspot.com
NDB's concurrency model is "lock-aware programming". You, as a programmer, decide whether you need to lock a row for writing or reading or whether you don't need a lock at all. Calvin serializes transactions for you, which is great, but you pay the price in terms of scalability (nowhere near NDB) and latency (nowhere near NDB). Spanner is a global OLTP DB, which is not comparable to NDB or Calvin.
But I'm a little confused by your comment: How is it possible to partition consensus without having more than leader? To me, the definition of "partitioned consensus" is that there is more than one consensus group, which means more than one leader.
Also, FYI, Calvin does not serialize transactions. It processes transactions in parallel. But it guarantees equivalence to a predetermined serial order. That distinction is important. As far as scalability, I discussed that in my previous post. Calvin doesn't have any scalability constraints that can be reached by known real-world workloads.
NDB has "lock-aware" programming - you don't get "global consensus". You decide, as a programmer, that this row could be accessed concurrently by another process, so you lock it, with either a read of write lock. Linearizability is easily implemented by acquiring a lock on a well-known row, but, of course, kills scalability.
In our Usenix FAST paper on HopsFS on Spotify's Hadoop workload, we had 1m ops/sec on HDFS, which was about 10m ops/sec on NDB. We ran out of hardware. There are workloads that big. [edited for clarity]
masterful understatement imo
In the context of this blog post, he specifically calls out partitioned consensus databases for requiring two wide-area round trips in order to run 2PC. However, we've seen multiple examples of partitioned databases (i.e. MDCC, TAPIR, Janus, and others) since the Spanner paper that can commit multi-partition transactions in a single wide-area round trip . Just as in the "unified consensus" approach, failures or concurrency may cause these systems to infrequently take multiple wide-area round trips to commit.
The blog post does a great job explaining the differences between "Calvin-like" systems and "Spanner-like" systems, but it falls short in convincing me that the "Calvin-like" architecture is fundamentally better, or makes better tradeoffs, than any partitioned architecture.
I implemented projects in GAE Datastore, though it's classic Paxos inside, it's clearly multi-partition database. For my workflow almost all of the frequqnt transactions were single-entity, so multi-partition Datastore worked fine with it.
EDIT: To be clear, when I say better I mean higher throughput and lower latency.