Cross shard transactions at 10M requests per second
blogs.dropbox.com
blogs.dropbox.com
> Edgestore data was already well-collocated, which meant that cross-shard transactions ended up being fairly rare in practice—only 5-10% of Edgestore transactions involve multiple shards.
What if performant cross-shard transactions is a red herring, and the thing that we should be looking more into is reliable automatic data colocation to avoid performing cross-shard transactions as much as possible? There's decent amount of academic research around this, with projects like SWORD [1] and Schism [2] that study shard load balancing as a problem of hypergraph partitioning. It seems like this might be worth incorporating into commercial distributed database projects.
Edgestore's API is set up to shepherd users into good collocation patterns by default, and a lot of work over the past year or two went into improving collocation and educating users about best practices. The collocation efforts were actually orthogonal to implementing cross-shard transactions, but they were obviously very beneficial.
For some reason this reminds me of something like the entity group concept in Google's Megastore [1].
For example, imagine that you have a distributed key-value store. You want to users (on different shards) to either both be able to see a piece of content or neither. You can achieve this by allocating a key on either shard (or a third shard), writing a reference to each user’s shards. If all of those writes succeed, you can write to the new key which was allocated. If any writes fail, you can bail and your data is consistent. Wrap the above logic in a nice library and your service will scale horizontally.
Whether that complexity is actually difficult to deal with is a whole other question, I could easily believe that for some applications it's not but apparently for Dropbox it's worth implementing TPC to avoid dealing with it.
Giving up some flexibility in the types of atomic operations you can do for a simpler architecture is sometimes worth it.
Or otherwise this does not work for entities that will be changed.
You're right that most systems probably don't need 2PC, which is why Edgestore didn't include it until now. As mentioned in the post, we finally felt that we had reached the right balance of tradeoffs to justify the API-level primitive.
If I look at Dropbox, i am sure features related to sharing folders between people/organizations cross-shards. You can't avoid them if you want to offer a fully featured product.
After a quick skim of 2PC in Edgestore (sorry, no time), it is unclear if there is a single transaction coordinator (TC) or not. I assume it is a single TC - that you can scale it to 10m trans/sec is impressive. The really hard part is have multiple TCs and to design protocols to coordinate their recovery after failure. Here is a good example in the open-source NDB (MySQL Cluster) system - https://drive.google.com/file/d/1gAYQPrWCTEhgxP8dQ8XLwMrwZPc...
I would be interested in a granular point-by-point comparison between this layer and native MySQL replication. In other words, if someone were to rewrite MySQL replication such that it did everything that this abstraction layer did (in addition to replication), what would it need to do?
Now (and this is strictly academic/theoretical), I'm curious what would be necessary to modify in the abstraction layer if the abstraction layer was to support a whole bunch of disparate SQL databases underneath it, i.e., Postgres, SQLite, SQL Server, Oracle, etc. (No, that wouldn't be practical, but it would be an interesting exercise to really learn where the gotchas might be where working with different SQL dialects...)
(I work at Google, but not on Vitess or YouTube)
I've had some opportunity to work with the Vitess (now PlanetScale) folks on some other databases-related work and it's a great community and product so I would definitely encourage people to check it out.
Twitter did it eight years ago with Gizzard. See http://highscalability.com/blog/2011/12/19/how-twitter-store... and https://blog.twitter.com/engineering/en_us/a/2010/introducin...
Disclaimer: I work on this database.
Amazon's Sable is built on top BerkeleyDB, acc to this: https://news.ycombinator.com/item?id=17595644