Notes on the Amazon Aurora Paper
blog.the-pans.com
blog.the-pans.com
That said, this can be solved by changing the level of abstraction. If you tightly coupled the execution and storage engine as a single module, which is not intrinsically dependent on the type of data model, it would allow high-performance multi-use within reasonable constraints. Albeit with a very different API than a storage engine. But that is not a thing that exists in open source AFAIK.
The loss from separating storage and execution may not matter for databases where operational efficiency is not paramount (e.g. because the scale is too small for it to matter) but many database engine designers are not going to take that use case limiting performance hit because the CapEx/OpEx implications are large.
The performance gains at the distributed level seem to rely on either parallelized compute (a-la MapReduce) or over-indexing (a-la ElasticSearch/GIN) rather than the sort of efficiency gains that were made at the level of optimizing between cpu/cache/memory/disk. I think modern problems regarding distributed databases are concerned with the time complexity not exploding as the size of the data does.
The distributed databases in open source are a bit of a biased sample as they tend to copy architectures from each other and are a trailing indicator of the state-of-the-art. For example, you virtually never see facilities for individual nodes to self-orchestrate parallel query execution (this is almost always centralized), which entirely changes how you would go about scaling out e.g. join operations or mixed workload. Building this facility typically requires a tightly coupled storage and execution engine, but open source systems are strongly biased toward decoupling these things for expediency. Open source database architecture is trapped in a local minima.
This is what we are doing with FaunaDB; Fauna's own native semantics are essentially a unified intermediate representation of a variety of high level query languages, forthcoming.
But by far the biggest advantage is operational. Your DB is now stateless! Nodes can die and be replaced in seconds, because they don't store any data; this makes designs like Bigtable's much more practical. You can quickly scale up the processing to deal with spikes without spending hours re-replicating. Operating the storage layer is still inherently challenging, but just has to be figured out once for all your DBs and other processing systems.
While performance can take a hit, this can be mitigated by local caching and good internal networking. It isn't able to meet the lowest levels of latency (sub ms), but in practice it works great for web applications.
Though later it went stateless.
TiDB has a similar architecture: https://www.pingcap.com/docs/architecture/
See more https://docs.datomic.com/on-prem/storage.html#storage-servic...
https://www.brentozar.com/archive/2019/01/how-azure-sql-db-h...
I sketched out a lot of architecture diagrams in that, showing how it differs from traditional RDBMS’s, and I tried to keep the post approachable for database folks from different platforms.
I hope they fix it
The comparison in the paper is against the synchronous-mirroring pattern underlying RDS Multi-AZ [1], which is the existing fully-managed, high-availability MySQL solution Amazon offers, so the direct comparison is appropriate for existing RDS Multi-AZ users considering a migration to Aurora.
[1] https://aws.amazon.com/blogs/database/amazon-rds-under-the-h...
Another thing I noticed with Aurora is the incremental cost of storage for Aurora is extremely cheap compared to the cost of storage in EBS or the cost of storage if you used instance storage in EC2. The Aurora storage is 0.1/GB but this is replicated 6 times which is more than EBS and EBS costs the same 0.1/GB. This also might be why no-one is going to build a similar system to Aurora. It will be hard to sell something like Aurora to cloud users because the storage costs are going to more than what EC2 charges.
For example, with something like "INSERT INTO my_table(field1, field2, field3, ...) SELECT f1, f2, f3... from root_table WHERE ...", I'm wondering how it will spread the load across nodes.
Aurora splits out 'database' nodes (the server instances you provision and pay for) from 'storage' nodes (a 'multi-tenant scale-out storage service' that automatically performs massively-parallel disk I/O in the background). Instead of MySQL writing various data to tablespaces, redo log, double-write buffer, and binary log, Aurora sends only the redo-log over the network to the storage service (in parallel to 6 nodes/3 AZs for durability).
No need for extra tablespace, double-write buffer, binary-log writes, or extra storage-layer mirroring, since durability is guaranteed as soon as a quorum of storage nodes receives the redo-log. The reduced write amplification results in 7.7x fewer network IOs per transaction at the 'database' layer for Aurora (vs standard MySQL running on EBS networked storage, in the benchmark described in the paper), and 46x fewer disk IOs at the 'storage' layer [1].
[1] https://www.allthingsdistributed.com/files/p1041-verbitski.p...
LOL
Though google has done this for all databases (at least their internal ones) (even bigtable). And some companies say they can do it for nvme & ram.
You just need real-good networking and some type of asic and removing kernel from the path.