Ensuring compatibility from a clean slate is hard, but to make MySQL distributed requires more than the current storage engine API provides. Thus, you would be likely be creating a fork and no longer getting the benefits of having an upstream[1].
The surgery you need to perform on the code would also reduce you from full compatibility (switching from pessimistic to optimistic locking means some statements no longer work). There are some great side benefits to a clean slate too: In TiDB we're able to use modern languages such as Go for the TiDB Server, and Rust for TiKV.
([1] Phrasing this question to upstream: I think they'd like to offer such functionality, but introducing this level of backward incompatibility is hard. They are also very concerned about performance regression bugs, which would be very hard to prevent with all the refactoring required.)
I doubt that the same will be done with PostgreSQL or MariaDB anytime soon. People are begging for a good and easy psql HA cluster setup for ages, and it is just not happening (yeah, they are getting closer and closer), while such DBs fill a need with good results.
For example, for local testing, in our Quick Start docs section, we have info for: mac/linux, involving just downloading and unpacking the release and you should be good to go; for local docker, you can download our control script; for k8s download our sample yml.
Finally, for non-local testing, we have a Deploy > Manual Deployment section, highlighting the 3 steps for downloading YugaByte, bringing up Masters (metadata nodes) and then Tservers (data nodes).
Note: I work at YugaByte.
With YugaByteDB: I needed to think about master and tserver nodes, to run them and connect them separately. There was no clear guide that had "run these three docker command on the three different nodes" and be done with it.
Note: I run everything is in simple docker images (and not compose, not swarm, not kubernates, just plain docker images).
Most other databases insist on multiple layers and dependencies which only shifts more work onto developers who have to install and maintain it.
> People are begging for a good and easy psql HA cluster setup for ages
I've been using Stolon in production for a bit, and while not a native solution- it works quite well.
To use a solution like Stolon would require changes in the core logic to dynamically re-configure whatever sharing plugin you use.
PG is planning the pluggable storage API for version 12.0, but these will only address user tables. This is certainly essential, but not sufficient for true distributed SQL. Some of the other pieces are: * Pluggable storage for system tables * Ability to create the initial set of system tables in a pluggable manner (initdb equivalent) * A number of other changes to the upper half of PG in order to make it "stateless" and scale-out, such as modifications to how certain operations are executed, which locks are held, what operations are pushed down to the underlying storage layer, etc. * Support for different types of indexes * Enhancing the optimizer to understand different storage types (in this case a distributed store)
In our current implementation, we have tried to make APIs for the above where ever possible. We need to explore how we can make runtime hooks for the other places. The plan is to see how to contribute these to the open source over the longer time frame so that we can create a self contained extension (though we are nowhere close to that today). Since our PG modifications code-base is open-sourced, the hope is that it would make contributing these changes back to PG easier.
At a high-level, the upper half of the Postgres DB is being largely reused. The lower-half, i.e. the distributed table storage layer, uses YugaByte's underlying core-- a transactional and distributed document-based storage engine.
For the DB to be scalable, the lower-half being distributed is necessary but NOT sufficient. The upper-half also needs to be extended to be made aware of other nodes executing DDL/DML statements and dealing with related concurrency while still allowing for linear scale. Also, making the optimizer aware of the distributed nature of table storage is the other major piece of work in the upper-half.
These changes required in the upper half is what makes the "100% pure extension" model a bit harder... but that's something we intend to explore jointly with the Postgres community.
CockroachDB gives you proper data replication and high-availability at a high density. Postgres doesnt, and requires an entirely separate cluster with fragile replication and manual cutover.