Umbra: an ACID-compliant database built for in-memory analytics speed
umbra-db.com
umbra-db.com
It does make me wonder what will be the next big leap in DB technology. Most of the NoSQL or distributed DB implementations have a bunch of limitations which make them impractical (or not worth the trade offs) for most applications, IMO. Distributed DBs are great until things go wrong, and then you have a nightmare on your hands. It's a lot easier to optimize simple relational DBs with caching layers, and adding read replicas scales quite effectively too.
The only somewhat recent new DB that comes to mind which had a really interesting model was RethinkDB, although it suffered from a variety of issues, including scale problems.
Anyway, these days I stick with Postgres for 99% of things, and mix in Redis where needed for key/value stuff.
However, I do think the key innovations are in building control planes around existing relational and NoSQL databses for scaling/sharding them across a set of resources to minimize cost while meeting performance and availability constraints.
RethinkDB, CockroachDB and FoundationDB are worth keeping an eye on.
Still, I can only be carefully optimistic.
Yeah, MySQL is ahead of the curve there...
Now, one approach is to just dismiss this use-case by pointing at DynamoDB and similar offerings. But if for some reason you can't use these hosted platforms, what do you use instead?
For search, ElasticSearch fortunately fits the bill, the "just keep adding boxes" concept works flawlessly, operating it is a breeze. But you probably don't want to use ElasticSearch as your primary datastore, so what do you use there? I had terrible experiences operating a sharded MongoDB cluster and my next attempt will be using something like ScyllaDB/Cassandra instead since operations seem to require much less work and planning. What other databases would offer that no-advance-planning scaling capability?
Somewhat unrelated, by I often wonder what one were to use for a sharded/distributed blob store that offers basic operations like "grep across all blobs" with different query-performance than a real-time search index like ElasticSearch. Would one have to use Hadoop or are there any alternatives which require little operational effort?
But you can still make this work. In our case, we'd ended up going the first route, and we ended up adding AWS EBS as the filesystem block store for our Postgres database, which is easy to resize dynamically without incurring downtime or other issues.
The downside of a vertical scaling approach is, well, you don't get HA "for free". You have to manually configure followers and standby nodes for the sole machine. You have to worry about your own replication. If failover and management is abstracted out for you, as in RDS, then it's easy to live with - otherwise very painful.
tl;dr go with managed DBs if you just need a DB and don't particularly have to optimise query performance outside of DB config.
Plus, a single server with many TB is a big single-point-of-failure, and if you want to scale it (vertically) you still have to take it down.
What kind of data and amounts do you frequently encounter that is a good fit for relational storage but difficult to fit in a single server? Honestly curious.
And you say update... Across keys, or some keys depend on others? Otherwise it sounds like a perfect fit for sharding with reduced constraints/requirements.
You can even use the Foundation Document Layer which is API compatible with Mongo.
It definitely takes some getting used to, but I think it's pretty fucking great, once you do.
CouchDB is re-architecting onto it, https://youtu.be/SjXyVZZFkBg
I think the biggest reason for a lack of noise about it, is the overhead of learning it is pretty high. You're not going to find folks writing their first "Nodejs + Express" applications using it. Additionally, you really have to know why most distributed databases suck ass to know why FoundationDB is so good.
Example: I have a cluster of three VMs running FDB on my home server, and over the past week I've accidentally hit the power switch three or four times. At no point did I have data loss in the cluster, the cluster "immediately" comes back up, and is ready to go. Adding machines to the cluster is unbelievably easy, especially if you have ever even tried grokking how to horizontally scale PSQL.
I'm close to releasing an Elixir-based Entity Layer which is pretty uhhh, ~~shitty~~ lacking, at this point feature wise, but it does make storing structured data a bit easier. I'm hoping it'll be more useful for helping people learn how to use FoundationDB than something folks are putting into production (though I'm dog-fooding it).
Yes, but it's also a wonderful way to corrupt a DB as soon as anything goes wrong in your storage fabric system.
And this type of corruption is generally not the recoverable one.
Why is this more likely to corrupt a DB than having it on local disks, when something goes wrong?
Local filesystems (ext4, xfs) have been designed to run over reliable internal bus and over spinning disk. They are (almost) able to stand most of the outage happening in this context ( powercut, corrupted block, missing flush ).
Put them over a non-reliable "normal" network, where you can get "savage" cable unplugged, faulty controller, paquet loss, out of order delivery, buggy middle box and you explode the number of scenario that can go wrong and will go wrong.
Block level I/O virtualization is amazingly useful, but (in my mind) it should be used with care...
I already heard of case in prod where the Master DB, the Slave DB and backup finished all in the same virtualized block-level storage.... try to guess what happened next.
Not even transparent pricing?
It’s admittedly not at all obvious from the repo name, but the following is the top level repo for Couchbase Server (which uses Google’s `repo` tool): https://github.com/couchbase/manifest
Geez, with all these limitations you may as well use another, proper open source database.
There was a really good post from Martin Fowler a while back that the popularity of "NoSQL" was really because it was "NoDBA" - app devs could sidestep the bottleneck of needing to get DBAs involved whenever you needed to persist an extra object field. While it's easy to abuse JSON storage in postgres, for things that are really just "opaque objects", vs. relational properties, appropriate use of JSON columns can save a ton of unnecessary overhead updating schemas.
In what way was it interesting? It was a document DBMS that supported MVCC.
They really shine for read heavy workflows that can tolerate a stale read every once in a while. If on top of that you have reasonable shard-ability you get near infinite scalability.
While that might cover a large portion of the database usage landscape, I'd hesitate to call it most. There's a reason OLTP was coined as an acronym--it's a pattern that comes up a fair bit.
Netezza, Informix, KDB and others will all well outperform open source dbs.
All open source dbs are absolute trash for time series an other olap type queries.
edit: after looking it up again, looks like that is still the case, and you have have to be fairly limited with cumulative aggregates if you want the keep performance. maybe someday, but as of now, still not very good.
I think Scidb, being a column-oriented dbms geared for multi-dimensional arrays (datacubes) is very interesting given current trends, and there are only a handful of similar dbs around, the other two that interest me are rasdaman and monetdb. I don't know if OmniScidb counts as a datacube db but it is also really interesting especially due to it's gpu and caching model.
As a sysadmin/ops type having to deal with monitoring timeseries db's are also something I like to keep an eye on. It used to be mostly rrdtool in this space, but now I am comparing prometheus and influxdb. Like OmniSci, another case of something that's not quite in the same db model but might even better a better solution for the space (metrics) is Apache Druid. (elasticsearch being another than can be massaged to fit as a quasi tsdb as well) I think there is some room to unify the monitoring/metrics and log storage arena into one space (usually they are separated which adds admin overhead) and right now I really like druid as a potential for this.
Another interesting application of timeseriesdb's that I have been keeping an eye on is in the quant/algo trading area. Most people have been using kdb+ there but many are looking for replacements and there are some really good conversations to be found about the kind of limitations they are hitting.
I'm just a sysadmin who likes to keep up with whats going on, and my db knowledge is limited, but I do have a process for narrowing my focus to the dbs I reference. It must be open source, bonus points to gpl or apache licenses. The language it is written in is important but not a deal breaker (very tired of so many java based dbs). I don't like it when they are tacked on top of an "older" tech (such as kairos on cassandra, timeseriesdb on top of postgres, opentdsb on top of hbase, kudu on hadoop, etc). Being either filesystem aware or agnostic can be nice features (playing well with ceph, lustre, etc) Not saying this is the sort of selection criteria others should use just giving some info on mine.
A few more interesting mentions: clickhouse, gnocchi, marketstore, Atlas (Netflix), opentick (on top of foundationdb).
Performance is good but comes at a massive usability cost. kdb+ brag about very succinct commands and even error messages but the small gain in efficiency from one less character in the error message is at the detriment of the user trying to work with it.
QuestDB (http://questdb.io) is an alternative to kdb as time-series data store (disclosure, i am one of the authors). We left trading to build it and show you can get the same speed without sacrificing usability. Users get better performance than kdb+, but can use SQL, ACID transactions. It is already in production with a few exchanges and HF, and it is open-source!
Hopefully, kdb users looking for a replacement find questdb helpful for their projects :-)
- storage engine written from scratch
- completely isolated read-only transactions and one read/write transaction concurrently with a single lock to guard the writer. Readers will never be blocked by the single read/write transaction and execute without any latches/locks.
- variable sized pages
- lightweight buffer management with a "kind of" pointer swizzling
- dropping the need for a write ahead log due to atomic switching of an UberPage
- rolling merkle hash tree of all nodes built during updates optionally
- ID-based diff-algorithm to determine differences between revisions taking the (secure) hashes optionally into account
- non-blocking REST-API, which also takes the hashes into account to throw an error if a subtree has been modified in the meantime concurrently during updates
- versioning through a huge persistent and durable, variable sized page tree using copy-on-write
- storing delta page-fragments using a patented sliding snapshot algorithm
- using a special trie, which is especially good for storing records sith numerical dense, monotonically increasing 64 Bit integer IDs. We make heavy use of bit shifting to calculate the path to fetch a record
- time or modification counter based auto commit
- versioned, user-defined secondary index structures
- a versioned path summary
- indexing every revision, such that a timestamp is only stored once in a RevisionRootPage. The resources stored in SirixDB are based on a huge, persistent (functional) and durable tree
- sophisticated time travel queries
As I'm spending a lot of my spare time on the project and would love to spend even more time, give it a try :-)Any help is more than welcome.
Kind regards Johannes
> - variable sized pages
> - lightweight buffer management with a "kind of" pointer swizzling
> - dropping the need for a write ahead log due to atomic switching of an UberPage
LMDB made those same design choices and is extremely fast/robust.
The most interesting bit is the use of variable size buffers (VSBs). The value of using VSBs is well known -- it improves cache and storage bandwidth efficiency -- but there are also reasons it is rarely seen in real-world architectures, and those issues are not really addressed here that I could find. Database companies have been researching this concept for decades. If one is unwilling to sacrifice absolute performance, and most database companies are not, the use of VSBs creates myriad devilish details and edge cases.
There are techniques that achieve high cache and storage bandwidth efficiency without VSBs (or their issues) but they are mostly incompatible with B+Tree style architectures like the above.
However, variable sized pages also allow page compression.
Can you give us some links to the mentioned issues and techniques that achieve high cache and storage bandwith efficiency without VSBs?
The alternative to VSBs is for each logical index node to comprise a dynamic set of independent fixed buffers, with each buffer having an independent I/O schedule. This enables excellent cache efficiency because 1) space is incrementally allocated and 2) the cache only contains parts of logical node that you actually use. References to the underlying buffers remain valid even if the index node is resized. Designs vary but 8 to 64 buffers per index node seems to be the anecdotal range. The obvious caveat is that storage structures that presume an index node is completely in buffer, such as ordered trees, don't work well. Since some newer database designs have no ordered trees at all under the hood, this is not necessarily a problem. There are fast access methods that work well in this model.
The main issue with VSBs is that it is difficult to keep multiple references to the page consistent, some of which may not even be in memory, since critical metadata is typically in the reference itself. A workaround is to only allow a single reference to a page, but this restriction has an adverse impact on some types of important architectural optimization. The abstract objective makes sense, but no one that has looked into it has come up with a VSB scheme that does not have these tradeoffs for typical design cases. That said, VSBs are sometimes used in specialized databases where storage utilization efficiency (but not necessarily cache efficiency or performance) is paramount, though designed a bit differently than Umbra.
The reason to use larger page sizes, in addition to being more computationally efficient, is that it gives better performance with cheaper solid-state storage -- storage costs matter a lot. The sweet spot for price-performance is inexpensive read-optimized flash, which works far better for mixed workloads than you might expect if your storage engine is optimized for it. Excellent database kernels won't see much boost from byte-addressable NVM and people using poor database architectures don't care enough about performance to pay for expensive storage hardware, so it is a bit of a No Man's Land.
> and subsequently allow the kernel to immediately reuse the associated physical memory. On Linux, this can be achieved by passing the MADV_DONTNEED flag to the the madvise system call.
Shouldn't this be MADV_FREE? This instantly reminded me of this classic Bryan Cantrill talk https://youtu.be/bg6-LVCHmGM?t=3529
Edit: It seems that the Linux behavior is relied upon? From later in the paper.
> Note that it is even legal for a page to be unloaded while the page con-tent is being read optimistically. This is possible since the virtual memory region reserved for a buffer frame always remains valid (see above), and read accesses to a memory region that was marked with the MADV_DONTNEED flag simply result in zero bytes. No additional physical memory is allocated in this case, as all such accesses are mapped to the same zero page
Well, that's a bold assumption as pg is speaking one of the richest sql dialects out there. And it also means it supports pg WAL protocol ?
The product is backed by solid research, so I suppose that there must be some powerful algorithms built-in, with a good coupling with hardware [1].
So the last question is how the code is made and tested, because good algorithms are not enough for a having a solid codebase. pg+(redis/memcached) is battle-tested.
Seems to use some common ideas with pg such as query jit compilation but mixes it with another approach.
> Umbra provides an efficient approach to user-defined functions.
possible in many languages using pg.
> Umbra features fully ACID-compliant transaction execution.
jepsen test maybe ?
Didn't harvest the clustering part neither.
[1] http://cidrdb.org/cidr2020/papers/p29-neumann-cidr20.pdf
I would really appreciate it.
The only bit I really understood was:
The system automatically parallelizes user functions
Now granted, I only understand how DBs work from a user-facing side so that might be a barrier here.Do you have any more information on this? I saw HyPer had been acquired by Tableau, and assumed it was a finished product they bought.
I'm perfectly willing to belief that they have no intention of selling. But that's really not a promise one can easily make. Even if you're capable of withstanding the allure of whatever large sum someone is offering, it's always possible to be faced with a choice of selling or shutting down, or selling or not being able to afford your spouse's/child's/own sudden healthcare needs.
That browser based query analyzer is cool.
Well that's impressive. Can I just drop this into my test suite and get a mega speed improvement? Could be worth it.
I especially like the super fast CSV scanning!
You will see a lot of people chasing the "Facebook wanna-be" kind of workloads.
I work with small/medium companies (or that are big, but with < 1TB of data). I bet 90% can't pass the first stages of data manipulation:
- Most(all?) rdbms have the same datatypes, mean: Use of nulls (bad) not algebraic types (sad), very unfriendly means to model data/business logic
- Use only SQL, that is impractical to anything bigger than basic queries. I work before with foxpro: I could do ALL THE APP with it, including GUIS, reports, etc. So to say SQL is disappointing is to say the less.
- All the engine stuff is a black box. That is great, until you wanna do your own index, store columnar data, save array or text or whatever you want, plug into the query executor and do your own stuff, etc.
You know, if you have JS and I tell you you can't code your own linked list, you will ditch that quickly. Sometimes, if you db engine allow to plug into the storage I could save my graphs for the only time I need it, instead of hack around putting it in tables or worse: Bring ANOTHER db engine to make my life hard.
Wait! Why?
All that stuff that some put in the middle-wares, model or controllers? In any other product you will reject the tool if you can't, but rdmbs FORCE to use anything else to finish it, despite the fact most run in the same box.
- Everyone add your own auth tables/logic, because the one implemented in rdbms is for a use case that not exist anymore. Then do it wrong, of course
- We are in the http world, but rdbms need something else for that.
- Import/Export data is still SAD.
- Import/Export data outside the few half-backed attempts (like csv, that a lot of time better you do it with python) is impossible in most. Using foreign adapters could work yet, you need to step out the black box, bring the adapter, compile it, install it, then use it. You need to become a C++ developer, despite you pretend to be a SQL one.
That is making thing worse because:
- Not exist a "rdbms" packager managers.
Look, how great if you just:
db install auth-jwt
db install csv-import
end. Like everyone else- Making of forms and reports. You can bet any company (even users!) will kill if their dbs allow to create reports and forms. Yep, alike access. Yep, that is because access is still a thing with the most weak db engine in town.
- You wanna allow to send emails, connect to external apis, call system services, etc.
But why? Is not that problematic? Well, if you JS (a language FOR THE BROWSER) allow it, why not your db? I live in that world before (foxpro) and it work GREAT.
---
A lot of this stuff is because the rdbms are look with a too narrow POV. Is crazy: People use half-finished NoSQL engines and are happy making his own query engine, yet talking about do else than SQL (or: plug into the query parser so i enhance it) will sound crazy to some.
Rdbms get leap-frog by NoSQL because until very late, get stuck in a mindset and use cases of the 80s.
Not exist ANYTHING that say you rdbms must be like all the others.
Broad the view is what, I think, rdbms need to do to get invigorated, and considering that are also performant, and very good, then will made to conquer the world!