8 ms·
Building CockroachDB on top of RocksDB
- polskibus 8y agoI noticed that RocksDB is used very often in OLTP scenarios. What's the OLAP equivalent of RocksDB in OLTP world? Apache Parquet? Apache Arrow? What would you use these days to create a high performance OLAP/OLHybridP engine ?
- deleted 8y ago[deleted]
- ryanworl 8y agoFor analytics workloads, your best bet is using compression techniques that let you do operations on the data without decompressing it. A good example is dictionary encoding a set of sorted string keys so you can preform prefix queries by doing a greater than and less than comparison on the integers instead of examining every string entirely. Once you’ve encoded the data into large enough blocks, you could use any storage engine and write the encoded blocks into it along with metadata for managing which blocks are a part of what tables and partitions of tables. You can also just use something like Parquet or ORC, but that’s not going to get you the best performance possible.
- polskibus 8y agoI know there are many techniques that used together give good performance (optimal memory layout, compression, vectorization, etc. etc.), however I'd like to use a package that does a lot of it, same what RocksDB (or SQLite) does for OLTP cases. Is there something like that? If not, what's out there that gives the best foundation for building OLAP functionalities on top of it?
- manigandham 8y agoThat's what Apache Arrow is, you had the right choice. That solves the processing component and you can use any number of on-disk formats like Parquet and ORC. And the hybrid of OLAP + OLTP is usually called HTAP.
- polskibus 8y agoIs it possible to insert new tuples to arrow model without rebuilding it from scratch?
- arjunnarayan 8y agoApache Arrow is your best bet, but it's still very much a new project without a lot of the things you're looking for.
- ryanworl 8y agoNo, I am not aware of any storage engine that provides that out of the box. The techniques are very tied into what your query processing engine can do and expects the data to look like. For example, do you materialize tuples immediately, or do you fully run it through your processing pipeline and not materialize until the end? Your storage engine and format needs to be at least somewhat involved in answer that question, because you need to know what data to read and when.
- arjunnarayan 8y agoUnfortunately most of the systems that build what you're describing are closed source (e.g. Snowflake, Microsoft SQL Server, Vertica, Teradata). There isn't an open-source project that does all of those things.
- georgewfraser 8y agoWhat about Presto?
- FridgeSeal 8y agoPresto is more of a distributed SQL solution: you run it on a cluster of nodes, point them at your storage later, it’s more optimised at querying very large datasets and it’s not built or tuned for high performance (in terms of latency or execution time).
- gianm 8y agoCheck out Druid [1], an open-source analytical database with tightly-coupled storage and processing engines designed for OLAP. In particular it implements a memory-mappable storage format, indexes, compression, late tuple materialization, and query engines that can operate directly on compressed data. There is a patch out to add vectorized processing as well, so you should expect to see that show up in a future release. Its storage format and processing engine aren't designed to be embedded in the same way as RocksDB and SQLite are, but you certainly could if you wanted to, since the code is fairly modular. Or you could use it as a standalone service as it was designed to be used. [1] http://druid.io/ http://druid.io/
- tristor 8y agoThere's also Clickhouse [1] which seems to scale much better than Druid, and has similar architectural decisions to make it somewhat general as a columnar store for OLAP uses. Cloudflare wrote an article in the past where they compared Clickhouse and Druid and they chose Clickhouse because they could get similar performance on the same workload with 9 nodes in Clickhouse which would require hundreds for Druid. They built all of the DNS analytics at CloudFlare on Clickhouse [2]. Disclosure: I work at Percona, and we've seen a lot of our customers make use of Clickhouse and have begun some of our own services work around it in Consulting. It's now a primary database talked about at our conferences, and we post about it regularly. [3] [1]: https://clickhouse.yandex/ https://clickhouse.yandex/ [2]: https://blog.cloudflare.com/how-cloudflare-analyzes-1m-dns-queries-per-second/ https://blog.cloudflare.com/how-cloudflare-analyzes-1m-dns-q... [3]: https://www.percona.com/blog/2018/10/01/clickhouse-two-years/ https://www.percona.com/blog/2018/10/01/clickhouse-two-years...
- polskibus 8y agoDoes ClickHouse support fine grained data security (for example role A gives access only to tuples with column X==123)?
- atombender 8y agoNo [1]. ClickHouse is a fairly low-level tool. If you need that kind of thing, you build an ACL-aware app on top of it. [1] https://clickhouse.yandex/docs/en/operations/access_rights/ https://clickhouse.yandex/docs/en/operations/access_rights/
- arjunnarayan 8y ago(author of the blog post here) I'd second ryanworl's comment that the rabbit hole goes much deeper than just storing things in a column oriented disk or in-memory format like Parquet or Arrow. That's just the first step. To get the best performance you have to have your data in an in-memory format that allows you to compress it efficiently, and then perform many relational operations on the compressed form itself. Another example is Run-length and delta encoding a sorted column of integers, and then building relational operators (e.g. a join) that operates directly on the compressed data. The best explanation for all the various techniques the go into the data structures and operator designs for OLAP workloads is the survey 'The Design and Implementation of Modern Column-Oriented Database Systems' by Abadi, Boncz, Harizopoulos, Idreos, and Madden: http://db.csail.mit.edu/pubs/abadi-column-stores.pdf http://db.csail.mit.edu/pubs/abadi-column-stores.pdf
- dominotw 8y ago> I noticed that RocksDB is used very often in OLTP scenarios. My experience with it has been most stream processing in kafka streams ect as local state store.
- jandrewrogers 8y agoThere is a practical engineering reason why the OLAP equivalent doesn't seem to exist. General purpose storage engines, and this applies to RocksDB, are like the C++ STL in that they provide good average performance across a wide range of common cases but are nowhere close to optimal if you have a well-defined type of data model and workload as your use case. You can always gain an integer factor increase in throughput by designing a less generalist implementation with a similar interface. As with the C++ STL, the limiting factor is the number of tunable parameters available i.e. the amount of internal architectural flexibility built into the implementation. OLTP storage engines are pretty simple, so a manageable number of behavioral parameters can usually get you within 3x of the throughput of a more targeted design, which is acceptable performance for most workloads that are not ingest-intensive. OLAP-ish storage engines, on the other hand, are at least an order of magnitude more complex to implement and have many more degrees of freedom depending on the expected data model and workload. There is a lot more data model and workload diversity in OLAP than OLTP, which makes implementing the effective internal architectural flexibility and set of tunable parameters that need to be maintained very unwieldy. If you limited yourself to the number of user-definable tuning and configuration parameters as an OLTP-oriented storage engine like RocksDB, the performance gap between a generalist implementation and a more targeted implementation will be more like 10-100x, which needless to say is huge. This makes the practical applicability of any "general purpose" OLAP storage engine that someone would want to use quite narrow, which diminishes the value of implementing a general purpose engine. This leads to the current reality that there is a zoo of specialist storage engines for OLAP-ish workloads -- graph, time-series, event processing, geospatial, classic DW, etc. Much more generalist OLAP storage engines that do several of these models could exist in theory but the bar for technical sophistication and complexity is much higher than for OLTP. Open source projects in particular tend to have a natural ceiling on the number of man-years invested to get an initial implementation of an architecture, which inherently limits the expressiveness of that architecture for software with this complexity.
- deleted 8y ago[deleted]
- StreamBright 8y agoS3 + ORC|Parquet + PrestoDB works very well
- ddorian43 8y agoSeems like few features: 1. sstables as different files 2. range delete (which is rare) compared to LMDB (which is faster & more efficient): https://symas.com/lmdb/technical/ https://symas.com/lmdb/technical/ Still would be nice to see how LMDB would fare in a complex distributed DBMS (most of them are in rocksdb-type libraries). But LMDB is supposed to stay small. So more features are in a fork: https://github.com/leo-yuriev/libmdbx https://github.com/leo-yuriev/libmdbx
- hyc_symas 8y agoLMDB is already used in distributed DBs - such as LDAP. OpenLDAP performance is orders of magnitude greater than any RDBMS or other distributed DB.
- perfmode 8y agoWhen a SQL implementation is built on a KV storage engine, how do tables, rows, and columns typically map to the underlying KV data model?
- arjunnarayan 8y agoExcellent question. There's a CockroachDB blog post about that: https://www.cockroachlabs.com/blog/sql-in-cockroachdb-mapping-table-data-to-key-value-storage/ https://www.cockroachlabs.com/blog/sql-in-cockroachdb-mappin...
- zuzun 8y agohttps://github.com/cockroachdb/cockroach/blob/master/docs/tech-notes/encoding.md https://github.com/cockroachdb/cockroach/blob/master/docs/te...
- perfmode 8y agoThanks, perfect!
- ryanworl 8y agoThe CockroachDB blog post on this topic is a good summary. There is an additional trick that isn't directly KV related, but is important in a distributed environment when using a KV storage engine. When defining a hierarchy of tables, such as customers -> orders -> order_line_items, you can make the primary key of the child tables contain the primary key of the parent table. e.g. (customer_id) for the customers table, (customer_id, order_id) for the orders table, then (customer_id, order_id, line_item_id) for the order_line_items. When this is stored on disk in a sorted format, it makes joins between these extremely cheap because the data will all be next to each other on disk. CockroachDB calls this "interleaved tables".
- gigatexal 8y agoInterleaved tables are best for 1:1 relationships.
- nanoseltzer 8y agoHas anyone done a study of how many potential users have been turned away by such a name that evinces disgust
- deleted 8y ago[deleted]
- the_duke 8y agoExcellent article, very informative. I just had to chuckle at this: > Non-engineers: in a computer, a move is always implemented as a copy followed by a delete Yeah, that's really gonna help a non-developer understand the article better...
- tyingq 8y agoIt's confusing altogether. For example, that's not how /bin/mv (usually) works.
- ulysses 8y agoUsually, that's because /bin/mv is just changing a link to the file, not moving the file itself. In cases where it's actually moving the file -- say across a file system boundary -- it does copy the file and then delete the old version.
- the_duke 8y agoI reckon it was meant in the context of compaction.
- tschellenbach 8y agoOur in-house DB at Stream also runs on top of RocksDB + Raft. Its amazing just how much faster it is than anything else out there (especially compared to cassandra). Instagram uses rocksdb as storage for Cassandra, Linkedin and pinterest use rocksdb. As soon as you have the time to build your own db using rocksdb you get really finegrained control over performance. https://stackshare.io/stream/stream-and-go-news-feeds-for-over-300-million-end-users https://stackshare.io/stream/stream-and-go-news-feeds-for-ov...
- stingraycharles 8y agoRocksdb is pretty good and we relied heavily on it at QuasarDB as well. Having said that, we are nowadays deploying more and more production setups with Levyx’ Helium, which scales better and directly integrates with the hardware.
- m0zg 8y agoGiven that Helium appears to be proprietary, what kind of perf benefit are we talking about here?
- stingraycharles 8y agoIn our testing, it’s multiple times faster, especially at scale. RocksDB’s compaction becomes a bottleneck fairly quickly when put under strain for extended periods of time. Helium performs much, much better at scale and doesn’t have compaction issues. It’s proprietary, but in my experience it’s money well spent. For the record, we were able to fully saturate a 4xNVMe with a 96 core server using Helium, while RocksDB achieved about 20% of the full NVMe capacity. As with all benchmarks, YMMV.
- ddorian43 8y agoDid you try LMDB ?
- jandrewrogers 8y ago
- dominotw 8y ago> If you surveyed most NewSQL databases today, most of them are built on top of an LSM, namely, RocksDB. Is this actually true? spark, foundationdb, memsql, nuodb , citus . I am not sure any of these are built on top of rocksdb. Which ones are actually built on lsm?
- nindalf 8y agoCassandra, MongoDB, BigTable, InfluxDB, LevelDB.
- dominotw 8y ago> Cassandra, MongoDB, BigTable, InfluxDB, LevelDB. None of these are NewSql[1](ACID and SQL) though. 1. https://en.wikipedia.org/wiki/NewSQL https://en.wikipedia.org/wiki/NewSQL
- manigandham 8y agoNone of those are newsql other than MemSQL, which is an OLAP system that uses a custom rowstore format and parquet for columnstores. In addition to CockroachDB there's also TiDB which runs on top of TiKV which uses RocksDB.
- dominotw 8y ago> None of those are newsql other than MemSQL Why aren't citus, nuodb 'newsql'? > there's also TiDB One more example doesn't qualify the statement "most are built on rocksdb". I wasn't saying there is only one newsql db built on rocksdb. Of the 14 examples listed here https://en.wikipedia.org/wiki/NewSQL https://en.wikipedia.org/wiki/NewSQL only 2 that you mentioned seem to be built on rocksdb.
- manigandham 8y agoI'm not disagreeing, RocksDB is not used by most. The statement in the blog post is not true.
- peterwwillis 8y agoRocksDB is a fork of LevelDB, which was [in]famous for its ease of corrupting data. Did Facebook ever do anything to ensure data wouldn't corrupt, or is that still a common thing operationally? (You find it more at larger scales) Here's an example of how data corruption can suck, with (example) Riak and LevelDB. The leveldb data would corrupt often, which would leave you in a predicament. Say you had 10 nodes with a 3 node replication factor, and the whole cluster is humming away at a decent clip. Now one node's leveldb corrupts, and you have to rebuild it. If you have a huge fuckoff dataset, this can take a while. Now another node goes down. Now only 1 node has the data you need, and 2 nodes are down - so now 8 nodes are doing the work of 10, and if you have any more failures, your data might be gone. Now add replication, which will suck performance and bandwidth away from the regular work. And because it would corrupt so easily & often, there needed to be hash trees to quickly identify what data was corrupt, and then you needed to fix it and rebuild your hash trees. This would also suck away performance. Finally, you can't just add new nodes while rebuilding, because the extra load makes the cluster fall over. And the more nodes, the higher the likelihood of failures.
- teacpde 8y agoCurious what is the cause of data corruption in leveldb?
- nickpsecurity 8y agoThat kind of concern is probably why FoundationDB built on and modified SQLite. Its reliability is already great.
- StreamBright 8y agoI never experienced this with several in production Riak clusters running for years. Can you explain how to reproduce or give a link to any public forum where this was discussed?
- peterwwillis 8y agoSure. Build about 10 classes of clusters of varying sizes, each with a dataset ranging from 100GB to a petabyte or more. Run them on shitty oversubscribed openstack clusters with a combination of ephemeral, Ceph, and SAN disks. Do replication to similar-ish clusters in different regions. Handle data for about 100 different applications that process so much data at such low latency that cloud-based databases aren't even an option. Keep adding nodes and storage to existing clusters over time. It turns out that really unstable hardware/networks like to expose bugs. It also wasn't discussed in public forums. We paid for support and even employed Riak developers, and still we hobbled on putting out fires. I'll bet other DBs go through the same crap and keep it quiet. Also, read the Riak documentation and you'll find the corruption recovery documentation among other hints at common failures and limitations.
- kureikain 8y agoIf someone love LevelDB/RocksDB but want to use a pure-Go implementation, I have good thing about this library: https://github.com/syndtr/goleveldb https://github.com/syndtr/goleveldb