Hydra: Column-Oriented Postgres
hydra.so
hydra.so
The problem with row based database structures is that the data between any two rows can vary in length, so if scanning across rows on disk for a column's value, you need metadata that points at the start of each column. That means the process of scanning across the data on disk involves a _lot_ of jumping around; less of a problem on SSDs than it was on spinning platters, but still work.
The alternate approach is to store columns together on disk. That means that when performing a query that filters based on values in a given column, you can scan sequential areas on disk, which is much much faster. Of course, if you now want to assemble rows based on your filter, you need to use the metadata about where related values are in each other column in order to assemble those rows, but hopefully you have a much smaller set of records that you need to do this over. (It's cheaper to assemble 10k matching rows across a set of columns than it is to jump across a billion unaligned columns to filter them in the first place.)
(If you have 100 rows of "1 Fizz" followed by 200 rows of "1 Buzz", followed by one "1 Fizz" you can store that as "100:1:Fizz;200:1:Buzz;1:1:Fizz".)
Then, a mere 30 bytes of IO + a quick sum tells you that you sold 101 Fizz and 200 Buzz from the widget factory this quarter. Maybe the first Fizz is how many you sold in the USA and the second Fizz is how many you sold in Canada, so it becomes easy to do by-region reporting.
However, unpacking the column that has all your individual transaction IDs and linking it to the individual single sale of one Fizz takes a lot of lookups. So untangling individual transactions becomes slow (though not impossible).
2) The operation of retrieving a full row of data is more expensive in a column based db than in a row based db, because in row based the data for the row is contiguous on disk. So row based is likely optimal if your data size works with that architecture. (And, to be clear, row based _can_ work with billions of rows, if you're thoughtful about what kind of querying you'll be doing against such tables, and how you maintain them.)
This is especially useful for reporting/analytical purposes where you want to aggregate over data from many records.
And do you run them in conjunction with your existing Postgres, install it as a plugin, or switch your entire db to something like Hydra?
As a rule of the thumb, if most of your queries are aggregations over few columns of many records (OLAP), column based would be a better choice - whereas if you need to access every column/field of a small number of records (OLTP), record/row based would be a better choice.
Do they also have a benefit similar to the compression that comes from record re-ordering in column databases?
It's not just that you need to read a smaller section of disk to scan the same amount of interesting data - it's that that data also takes up considerably less space than it otherwise would.
Yes, it depends on the data, but we've observed customers with a 10X data compression over Postgres (heap) tables although ~5-7X is more common.
In columnar databases a batch of column data is stored sequentially followed by the next column and so on. So if you need to sum a column you can read it sequentially very quickly and sum it up without needing to touch the data of any other column.
Columnar is also better at filtering and grouping for similar reasons.
As a general rule if your workload looks like “select sum(x) where y group by X” a columnar db will be a great fit.
The trade off is that columnar databases don’t like being updated and work best when you only insert data. This is because when you update existing data you now need to rewrite several pages on disk because the columns are all over the place.
Because of this there’s very few areas where a pure columnar DB makes sense. They’re great for timeseries metrics like sensor data , click event logs etc. places where you are just recording a fact and don’t need to update past events.
This is how many aggregate functions can run real fast.
To take advantage of that, data has to be structured a certain way: In a columnar way (which is then be vectorized and pushed to SIMD).
1. No indexes or primary keys
2. Insert and delete only, no updates
3. Automatically normalized behind the scenes leading to data compression in the 10x-100x range
4. Very fast analytical queries
BigQuery, RedShift, etc are columnar databases (for example).
You can convert a table in SQL Server to "Clustered ColumnStore" storage, it'll become compressed, but everything will mostly work the same.
The implementation in SQL Server uses a row "delta store" to keep recent changes. Once a certain threshold is reached, the delta store is converted to columnar format and merged with the underlying table.
[#001][John Smith][25];[#002][Jane Doe][32];[#003][Kane Citizen][62];[...
Columnar databases store the data as a "structure of arrays", keeping each type of value together:
[#001][#002][#003][...;[John Smith][Jane Doe][Kane Citizen][...;[25][32][62][...
With row-based storage, the database can look up everything it needs to show a form to the user with a single disk I/O operation: just read a block of data starting at the location where the record row is.
With column-based storage, the database has to do an I/O per column, which seems worse... but these days on SSD storage the overhead is negligible when displaying one form at a time.
Where columnar wins massively is reporting workloads. If you want to know the average age of employees, then with row storage you pretty much have to read everything in, including their names and ID numbers, even though you don't need that data. This is because all disk storage comes as "block devices" that read a minimum of 512 bytes, but typically more like 64KB at a time. You get more than you want, and then you have to throw it away. With column storage, you can read just the age data to compute the average age.
With column storage there are also compression techniques that work a lot better than with row storage. Run-length encoding, lookup tables, etc... provide higher compression ratios because the data is much more uniform for much larger chunks.
The typical end-result is 100x less I/O to run the same query.
In row oriented, each cell for each row is next to each other in memory. Think traditional CSV files, each row is a line in the file. To sum a column you have scan the whole file and all data for all columns.
It's like a database engine that uses grep under the hood.
Data is in flat files. Queries are executed by scanning those files (as quickly as possible).
For aggregate queries (which is what these things are used for), substitute awk for grep.
https://benchmark.clickhouse.com/#eyJzeXN0ZW0iOnsiQXRoZW5hIC...
shared_buffers = 128MB
max_parallel_workers = 8
max_parallel_workers_per_gather = 2
work_mem = 4MB
"Postgres (tuned)" uses: shared_buffers = 8GB (64x the default)
max_parallel_workers = 16
max_parallel_workers_per_gather = 8
work_mem = 4MB
"Hydra" uses: shared_buffers = 8GB
max_parallel_workers = 16
max_worker_processes = 32 (default 8)
work_mem = 64MB
The default settings are not realistic when running on a c6a.4xlarge instance, and the "Postgres (tuned)" settings more closely match what "Hydra" is configured to use.I consider it dishonest because I do not believe a team working this close to Postgres internals wouldn't know this.
very somehow, since Hydra essentially is allowed to use 4 times more CPU cores.
basically larger buffers and more concurrent workers.
* Citus (cstore_fwd) extension
* Azure CosmosDB for PostgreSQL
* GCP AlloyDB for PostgreSQL
* AWS Redshift (but very old Postgres base)
Fwiw, the main use case we had in mind when developing the Citus columnar table access method was compression of old data in a time-partitioned table. Citus is commonly used for real-time analytics on time series data. That involves more materialization than typical OLAP reporting use cases. However, you might still want to keep the source data in the database for ad-hoc queries or future materialization and it's preferable to keep the source data in compressed form.
Redshift is specifically architected for ad-hoc OLAP queries. AlloyDB I'm not sure.
https://www.cockroachlabs.com/blog/vectorized-hash-joiner/
https://www.cockroachlabs.com/blog/vectorizing-the-merge-joi...
...it seems the distinction here is that the vectorization is only present in the execution layer and not the storage layer also. I would guess that from a storage perspective, even with column families in play, everything is being streamed out of sorted a LSM engine regardless. So there isn't additionally some highly-tuned buffer pool serving up batches of compressed column files etc.
I am wondering why they are saying it is not for OLAP workload..
But I should look at TiDB, they looks like interesting and relatively mature project.
Another example is cassandra is not column oriented.