4 ms·
This is really hard to answer in a short comment, but I'll try to highlight a few things: - We use compute/storage separation as in Snowflake and Aurora. This
by mbravenboer 5y ago
This is really hard to answer in a short comment, but I'll try to highlight a few things:
- We use compute/storage separation as in Snowflake and Aurora. This means that data is stored in object storage (for example S3) and computing resources are transient. This is the de facto standard now for new database management systems. Systems that are not designed for object storage are gradually going to be less relevant. The Snowflake papers are a great resource on this ( http://pages.cs.wisc.edu/~yxy/cs839-s20/papers/snowflake.pdf http://pages.cs.wisc.edu/~yxy/cs839-s20/papers/snowflake.pdf ) . It's also a recurring topic in the CMU database talks. Our memory management is based on LeanStore and Umbra.
- The database is entirely versioned with an immutable data structure, so read-only scaling does not require locking: you can provision new compute instances that work on a snapshot of the database that is 100% guaranteed to be consistent. This means that you can run an analytical job on your production database with a 100% guarantee that the performance of your primary system is not impacted. For concurrent writes we aim to use incremental maintenance techniques.
- All our data is indexed so that the WCOJ algorithms can do their job. Indexes maintenance can be expensive, so we have write-optimized data structures for that (B-epsilon trees - http://supertech.csail.mit.edu/papers/BenderFaJa15.pdf http://supertech.csail.mit.edu/papers/BenderFaJa15.pdf ). The performance of those is really good, and they can be compacted in the background. The combination of our graph normal form (narrow relations) and WCOJ addresses the index selection problem in a novel way.
- We use new compiler architecture ideas to make everything live and demand-driven. Our framework for this is Salsa, which is open source ( https://www.youtube.com/watch?v=0uzrH2Ee494 https://www.youtube.com/watch?v=0uzrH2Ee494 and https://github.com/RelationalAI-oss/Salsa.jl https://github.com/RelationalAI-oss/Salsa.jl ).
- WCOJ are our work horse join algorithm. We have different query evaluation techniques, including JIT compilation and vectorization. For queries that do not require worst-case optimal joins the plans are similar to more conventional systems due to optimizations applied to the plan (talk: https://www.youtube.com/watch?v=C_mBEq_o4HE https://www.youtube.com/watch?v=C_mBEq_o4HE )
Graph workloads have a few interesting properties: 1) Analytical queries are often recursive. Systems like Neo4J do not really address this in general because the query language does not support general recursion. Instead, it offers bindings to algorithms implemented in Java. 2) Queries often do joins across many relations. 3) There is lots of skew.
(2) and (3) are handled by WCOJ. For (1) we use algorithms similar differential dataflow. Differential Dataflow already has great proof points for graph analytical queries (mainly created by Frank McSherry)
In the end, of course all that matters are the actual benchmarks. We are developed those for standard relational workloads (like TPC-H) and graph workloads (for example LDBC SNB, graph analytical workloads). It's a little too early for us to claim results at scale, but we have isolated proof points that the design works.