The design and implementation of modern column-oriented database systems
blog.acolyer.org
blog.acolyer.org
I also want to note that time-series databases are basically obsolete at this point because time-series data is very well handled by these OLAP column-stores. Create a table with time as a primary or sort key and you'll get fast queries with full SQL and joins.
There is no one-database-fits-all.
A database that performs well on TPC-H SF 100 will not necessarily perform well on TPC-H SF 10000.
One that performs well on TPC-H may be completely useless at TPC-DS.
So, no.
There's nothing stopping you from putting other data in there that doesnt have a time component at all, but it's designed to primarily work in the fintech space so has special performance optimizations and features for working with time and money.
Be careful with this. You can easily end up with pathologically bad performance, resource usage, contention etc. if most of your requests end up being directed to the same workers or shards in a distributed system. Partitioning by time can cause this because most queries tend to be for recent time windows.
Time can be a decent primary key withing a shard, but is normally a really bad choice for any key that's used to partition data between shards. Make sure you understand how your system partitions data so you don't fall into this.
Apologies if that came across as pedantic - you probably already know this, but hopefully some other readers can be saved from learning this through painful experience.
In a column store, you would apply the sort/cluster key to your time column, which would partition the table within nodes, not across nodes. This will speed up queries that specify a narrow time range, and will not cause any of the pathological behavior you describe.
Some column stores allow you to chhose a distribution key, which specifies how data is distributed across nodes. It can be used to optimize joins but it’s dangerous—it can cause exactly the problem you describe.
Also, time-series database implemented as you suggest will work in some cases, and better than some might expect if you are clever, with the caveat that it has significant weaknesses and scalability limitations. Many real-world time-series data models and workloads would run poorly on a database organized this way. Sorting or sharding on time in the data organization is usually a mistake in my opinion. Some high-performance time-series databases don't use time as a sort/shard key at all, instead leveraging the approximate partial order of the underlying data sources directly. This often offers orders of magnitude better write performance with little or no loss of query performance in practice and it generalizes better.
This is all a very complex but fantastically interesting technical design area. I firmly believe that scalable high-throughput mixed-workload analytical databases are possible, but they largely don't exist. Yet.
Yes, time as a secondary sort key works great. We have trillion row time-series tables with sub-second query times.
First, online insertion (parsing, processing, indexing, disk storage) of millions of records/second should be per node, basically saturating the ingress path of each node's 10/25-GbE interface, and scale-out additively. You'll need the aggregate insertion throughput of several nodes for many live sensor data sources, never mind trying to load several petabytes of historical data in a reasonable amount of time. Second, by implication, the workload is far too large for an in-memory database -- those online indexes you are dynamically constructing need to be on disk and distributed. Indexing several columns through disk storage at wire speed under these parameters is not something either MemSQL or SQL Server (or any database really) is designed to do, which was my point. Being in-memory often has few implications for query performance in well-designed parallel database kernels; on modern hardware, storage I/O is rarely the throughput bottleneck for mixed workloads if you are doing everything else optimally.
The core problem is that online disk-backed secondary index construction in a distributed environment that can use the full bandwidth of the switch fabric isn't feasible. A modern storage engine and scheduler will get you the raw throughput in most cases but that is a separate question from how do you support low-latency queries at the same time. Like I alluded to with time-series databases, the way you avoid the low throughput of index construction is by not constructing indexes -- other, sometimes exotic, selectivity mechanism exist that have orders of magnitude lower construction cost than a classic index while providing relatively efficient search. Even if those search mechanisms are slightly less selective than a traditional index, you make up for it with the intrinsically higher throughput of those mechanisms.
re: time-series, I don't doubt that you can do sub-second queries on a time-series table with a trillion rows -- I've done it with tens of trillions of rows in a single table. Time-series is one of the easiest data models to scale if you have a modicum of cleverness. The challenge is ensuring this query performance while inserting tens of millions of records per second into the index at the same time and/or you need to index more dimensions than a timestamp.
Log(Graph): A Near-Optimal High-Performance Graph Representation [pdf]
https://people.csail.mit.edu/jshun/papers/loggraph.pdf
For something more database specific, check out the new RedisGraph module [1] implemented with GraphBLAS [2] by the venerable Tim Davis [3] himself, who as you know implements the underlying sparse matrix algos used in everything from MATLAB to Google Maps...
RedisGraph in the Language of Linear Algebra with GraphBLAS, co-presented by RedisLabs and Tim Davis [video]
https://www.youtube.com/watch?v=xnez6tloNSQ
GraphBLAS is not Redis-specific, but the Redis data structures and new modular design made it an ideal candidate for the first popular database implementation.
This is significant because GraphBLAS is more than 10 years in the making, the culmination of the initial D4M matrix model design by Jeremy Kepner [4] and his team at MIT Lincoln Laboratory Supercomputing Center.
And for the last ~5 years or so the software model has been designed in collaboration with hardware teams at Intel, NVIDIA, IBM, and the labs [5] to make chips and architectures optimized for these new matrix models and capable of exascale.
It all came together this summer with the release of GraphBLAS 1.0 [6]. Now that GPU and TPU accelerators are populating the data centers and linear algebra has come en vogue, hopefully enough software engineers will be ready with the background understanding and a working mental model for the underpinning architectural paradigm shift [7].
[1] RedisGraph https://oss.redislabs.com/redisgraph/
[2] GraphBLAS http://graphblas.org
[3] Tim Davis http://faculty.cse.tamu.edu/davis/
[4] Jeremy Kepner http://www.mit.edu/~kepner/
[5] GraphBLAS: Building Blocks For High Performance Graph Analytics https://crd.lbl.gov/news-and-publications/news/2017/graphbla...
[6] Graph algorithms via SuiteSparse:GraphBLAS: triangle counting and K-truss [pdf] http://faculty.cse.tamu.edu/davis/GraphBLAS/HPEC18/Davis_HPE...
[7] David Patterson Says It’s Time for New Computer Architectures and Languages https://news.ycombinator.com/item?id=18009581
> Open Source Column-Oriented Databases: There are three available open source column databases, all were based on works of research groups that later saw commercial spinoffs. C-Store produced vertica, MonetDB spawned Vectorwise and LucidDB was DynamoBI. Each project has stagnated as the team around them moved on to commercial endeavours. Only MonetDB appears to be still actively developed, it's also the one that seems most feature complete as we'll see later.
Is the above an accurate characterization of OSS column-oriented databases?
There is also MapD which uses GPUs, https://github.com/mapd/mapd-core
> MapD Core is an in-memory, column store, SQL relational database that was designed from the ground up to run on GPUs.
Cassandra is a https://en.wikipedia.org/wiki/Wide_column_store not column oriented: https://en.wikipedia.org/wiki/List_of_column-oriented_DBMSes
Cassandra is a (badly named) wide-column database, which is more accurately described as an advanced key/value store, specifically a sorted hash map with multiple levels of values, which can also be part of the primary and secondary keys.
- Appending cells to a column
- Removing a range of cells from a column
- Searching for data in a column (exact search or complicated text search).
- Performing a map over a column (applies user defined function to each cell, producing a new column or makes changes in-place)
- Filtering a column using a user-defined function.
- Combining two or more columns through some user-defined function
- Sorting
- more ?
+ "Appending cells to a column" ... when writing data to a "table" the columnar database goes through a lot of effort to compress each column of data into separate files on disk -- the benefits of compression come from regarding many "rows" of data, per column, at once, so appends probably go to transient areas, and then in the background recent appends are merged-and-compressed into files on disk ... data in all columns in the columnar DBs usually are stored sorted per some timestamp key, and that data doesn't always come in perfect order (e.g., there are different latencies for log records from the same timestamp arriving from different hosts / devices / applications in an SIEM system) -- so appends will be fast to keep up with the huge ingestion flow, but a lot of painstaking expensive compression is done while merging, and re-merging data -- most columnar DBs systems will segment data based on buckets of the timestamp range (and perhaps the more dense the table time-wise, the narrower are the buckets)
+ "Removing a range of cells from a column" ... removing a range of timestamp-ordered or even worse a filtered set of cells will be very expensive -- for all columns in the range to be removed, all deleted items will eventually need to be removed (e.g., 100 and 200 rows in two different timestamp buckets, per column, of 1,000,000 rows of data per bucket) -- reading, uncompressing, and filtering out removed columns will be fast -- recompression and re-writes will be slow ... there may be ways these DBs use side-data to instead, for a period before "vacuum" or in perpetuity, mark rows as deleted rather actually removing and recompressing the data -- these DBs are usually meant for write-once-and-never-delete, read-many-times situations
+ "Searching for data in a column" -- the main benefit from columnar DBs is that, once the data is available in compressed form on disk ... a query needs to read from disk only the specific columns involved in the search filter/terms, plus those needed to generate the output payload -- a typical query might surface just 5 of 80 columns, so compared to a traditional DB that reads all the data for each row (if data strays outside indexed columns), a columnar search will quickly fetch only the already-compressed data for a small number of columns off disk. queries usually involve scanning a huge number of rows in a timerange often expanding multiple timestamp buckets -- for each column, if the columns were stored uncompressed, the disk latency and bandwidth for reading the huge number of row values per column would dwarf the time taken in CPU/memory for processing the data -- with a columnar-DB with compressed columns, each disk read (the slow part) yields a huge number of easily-and-quickly-decompressed row values from that column. the compression is slow, but decompression is very fast ... so many rows of data, per column, can be pumped out to filtering/aggregation portions of the engine, much faster than in a traditional DB (I worked at a place that got 40x compression over traditional DBs)
+ "Performing a map over a column" -- the DB engine can process many rows across the columns of data involved quickly enough -- filters based on value are easy, functions can be applied to columns of data, or data from multiple columns can be joined -- as long as computations involve only data per-row, is usually fast -- it's additional aggregation, and multi-level query processing (results of one query, perhaps aggregated, feed into another), that involve careful planning of a query, especially in a distributed cluster where aggregation "keys" may require retrieval of rows from shards on different hosts
+ "Filtering a column using a user-defined function" ... fast, just like prior point
+ "Combining two or more columns through some user-defined function" ... fairly fast
+ "Sorting" ... here it would get tricky, because to sort a bunch of raw records without filtering takes a huge number of resources ... unless your sort is already on the timestamp field that's the organizing sort field for the columnar DB ... most queries won't involve retrieving 5 hours of data, sorting the raw data by some field other than "main timestamp", and returning results ... usually data is fetched, some aggregation is done in the query (perhaps multiple levels), and sorting is done on the vastly reduced (several orders of magnitude) data -- this kind of sorting is much easier
+ ...
+ "Joins" ... joins between two huge datasets on the raw contents of each table will be real expensive ... so expensive that the columnar DB vendor I worked at (starting in 2000) chose to not implement joins other than joining massive-DB-table-data with vastly smaller lookup-tables -- the value proposition for us was not in doing the joins, it was vastly reduced, reliable storage of huge numbers of log records, resulting in cheaper compression and very fast querying in a distributed cluster
> The ultimate performance boost from compression comes when operators can act directly on compressed values without needing to decompress. Consider a sum operator and RLE encoded values. It suffices to multiply the value by the run length.
It sounds very impressive and I can see in theory how it could benefit COUNT(), COUNT(DISTINCT ), SUM(), AVG(), MIN(), and MAX()... in other words, pretty much all of the common aggregation functions. For example, COUNT(DISTINCT *), MIN(), and MAX() can literally just ignore run length as those are known duplicates, while the others simply need to use the run length as a factor.
Does anyone know if this specific optimization is actually implemented in any real-world database?
They keep what are called "zone maps", which are pre-computed statistics for groups of rows that contain those same aggregate functions. This allows for those types of analytics functions, but also for pruning of data to read at all during regular table scans further in the query execution.
explain select country, count(*) from tablename group by 1
...
+-GROUPBY PIPELINED (GLOBAL RESEGMENT GROUPS) [Cost: 79, Rows: 1] (PATH ID: 1)
| Aggregates: count(*)
| Group By: tablename.country
| Execute on: All Nodes
| +---> STORAGE ACCESS for tablename [Cost: 78, Rows: 14M (1 RLE)] (PATH ID: 2)
...
Note the row count at the top level.This type of thing only works when the group column is the first in the projection sort order (or for multiple columns, a prefix of it). Generally GROUPBY PIPELINED indicates that the sort order is being used (including RLE).
That said, some traditional row-store databases like SQL Server support columnar indices, which give row-oriented databases similar performance characteristics to columnar. The downside is lower write speed (the columnar indices have to be maintained, and the writes are still done row-wise).
All indices are derived views of the original dataset that are optimized with certain characteristics.
A column-oriented table can also have regular b-tree indexes as well. Perhaps it's unfortunate that they used indexes as the main interface for managing these table types instead of explicit language.